--- /dev/null
+// Licensed to the .NET Foundation under one or more agreements.
+// The .NET Foundation licenses this file to you under the MIT License.
+// See the LICENSE file in the project root for more information.
+
+using System.Reactive;
+using System.Threading.Tasks;
+
+namespace System
+{
+ public static class AsyncObservableExtensions
+ {
+ public static ValueTask<IAsyncDisposable> SubscribeAsync<T>(this IAsyncObservable<T> source, Func<T, ValueTask> onNextAsync)
+ {
+ if (source == null)
+ throw new ArgumentNullException(nameof(source));
+ if (onNextAsync == null)
+ throw new ArgumentNullException(nameof(onNextAsync));
+
+ return source.SubscribeAsync(new AsyncObserver<T>(onNextAsync, ex => new ValueTask(Task.FromException(ex)), () => default));
+ }
+
+ public static ValueTask<IAsyncDisposable> SubscribeAsync<T>(this IAsyncObservable<T> source, Func<T, ValueTask> onNextAsync, Func<Exception, ValueTask> onErrorAsync)
+ {
+ if (source == null)
+ throw new ArgumentNullException(nameof(source));
+ if (onNextAsync == null)
+ throw new ArgumentNullException(nameof(onNextAsync));
+ if (onErrorAsync == null)
+ throw new ArgumentNullException(nameof(onErrorAsync));
+
+ return source.SubscribeAsync(new AsyncObserver<T>(onNextAsync, onErrorAsync, () => default));
+ }
+
+ public static ValueTask<IAsyncDisposable> SubscribeAsync<T>(this IAsyncObservable<T> source, Func<T, ValueTask> onNextAsync, Func<ValueTask> onCompletedAsync)
+ {
+ if (source == null)
+ throw new ArgumentNullException(nameof(source));
+ if (onNextAsync == null)
+ throw new ArgumentNullException(nameof(onNextAsync));
+ if (onCompletedAsync == null)
+ throw new ArgumentNullException(nameof(onCompletedAsync));
+
+ return source.SubscribeAsync(new AsyncObserver<T>(onNextAsync, ex => new ValueTask(Task.FromException(ex)), onCompletedAsync));
+ }
+
+ public static ValueTask<IAsyncDisposable> SubscribeAsync<T>(this IAsyncObservable<T> source, Func<T, ValueTask> onNextAsync, Func<Exception, ValueTask> onErrorAsync, Func<ValueTask> onCompletedAsync)
+ {
+ if (source == null)
+ throw new ArgumentNullException(nameof(source));
+ if (onNextAsync == null)
+ throw new ArgumentNullException(nameof(onNextAsync));
+ if (onErrorAsync == null)
+ throw new ArgumentNullException(nameof(onErrorAsync));
+ if (onCompletedAsync == null)
+ throw new ArgumentNullException(nameof(onCompletedAsync));
+
+ return source.SubscribeAsync(new AsyncObserver<T>(onNextAsync, onErrorAsync, onCompletedAsync));
+ }
+
+ public static ValueTask<IAsyncDisposable> SubscribeAsync<T>(this IAsyncObservable<T> source, Action<T> onNext)
+ {
+ if (source == null)
+ throw new ArgumentNullException(nameof(source));
+ if (onNext == null)
+ throw new ArgumentNullException(nameof(onNext));
+
+ return source.SubscribeAsync(new AsyncObserver<T>(x => { onNext(x); return default; }, ex => new ValueTask(Task.FromException(ex)), () => default));
+ }
+
+ public static ValueTask<IAsyncDisposable> SubscribeAsync<T>(this IAsyncObservable<T> source, Action<T> onNext, Action<Exception> onError)
+ {
+ if (source == null)
+ throw new ArgumentNullException(nameof(source));
+ if (onNext == null)
+ throw new ArgumentNullException(nameof(onNext));
+ if (onError == null)
+ throw new ArgumentNullException(nameof(onError));
+
+ return source.SubscribeAsync(new AsyncObserver<T>(x => { onNext(x); return default; }, ex => { onError(ex); return default; }, () => default));
+ }
+
+ public static ValueTask<IAsyncDisposable> SubscribeAsync<T>(this IAsyncObservable<T> source, Action<T> onNext, Action onCompleted)
+ {
+ if (source == null)
+ throw new ArgumentNullException(nameof(source));
+ if (onNext == null)
+ throw new ArgumentNullException(nameof(onNext));
+ if (onCompleted == null)
+ throw new ArgumentNullException(nameof(onCompleted));
+
+ return source.SubscribeAsync(new AsyncObserver<T>(x => { onNext(x); return default; }, ex => new ValueTask(Task.FromException(ex)), () => { onCompleted(); return default; }));
+ }
+
+ public static ValueTask<IAsyncDisposable> SubscribeAsync<T>(this IAsyncObservable<T> source, Action<T> onNext, Action<Exception> onError, Action onCompleted)
+ {
+ if (source == null)
+ throw new ArgumentNullException(nameof(source));
+ if (onNext == null)
+ throw new ArgumentNullException(nameof(onNext));
+ if (onError == null)
+ throw new ArgumentNullException(nameof(onError));
+ if (onCompleted == null)
+ throw new ArgumentNullException(nameof(onCompleted));
+
+ return source.SubscribeAsync(new AsyncObserver<T>(x => { onNext(x); return default; }, ex => { onError(ex); return default; }, () => { onCompleted(); return default; }));
+ }
+ }
+}
\ No newline at end of file
--- /dev/null
+// Licensed to the .NET Foundation under one or more agreements.
+// The .NET Foundation licenses this file to you under the MIT License.
+// See the LICENSE file in the project root for more information.
+
+using System.Threading.Tasks;
+
+namespace System
+{
+ public interface IAsyncObservable<out T>
+ {
+ ValueTask<IAsyncDisposable> SubscribeAsync(IAsyncObserver<T> observer);
+ }
+}
\ No newline at end of file
--- /dev/null
+// Licensed to the .NET Foundation under one or more agreements.
+// The .NET Foundation licenses this file to you under the MIT License.
+// See the LICENSE file in the project root for more information.
+
+using System.Threading.Tasks;
+
+namespace System
+{
+ public interface IAsyncObserver<in T>
+ {
+ ValueTask OnNextAsync(T value);
+ ValueTask OnErrorAsync(Exception error);
+ ValueTask OnCompletedAsync();
+ }
+}
\ No newline at end of file
--- /dev/null
+<Project Sdk="Microsoft.NET.Sdk">
+
+ <PropertyGroup>
+ <TargetFramework>net5.0</TargetFramework>
+ <RootNamespace>System</RootNamespace>
+ </PropertyGroup>
+
+</Project>
--- /dev/null
+// Licensed to the .NET Foundation under one or more agreements.
+// The .NET Foundation licenses this file to you under the MIT License.
+// See the LICENSE file in the project root for more information.
+
+using System.Threading.Tasks;
+
+namespace System.Reactive
+{
+ public class AsyncObservable<T> : AsyncObservableBase<T>
+ {
+ private readonly Func<IAsyncObserver<T>, ValueTask<IAsyncDisposable>> _subscribeAsync;
+
+ public AsyncObservable(Func<IAsyncObserver<T>, ValueTask<IAsyncDisposable>> subscribeAsync)
+ {
+ _subscribeAsync = subscribeAsync ?? throw new ArgumentNullException(nameof(subscribeAsync));
+ }
+
+ protected override ValueTask<IAsyncDisposable> SubscribeAsyncCore(IAsyncObserver<T> observer)
+ {
+ if (observer == null)
+ throw new ArgumentNullException(nameof(observer));
+
+ return _subscribeAsync(observer);
+ }
+ }
+}
\ No newline at end of file
--- /dev/null
+// Licensed to the .NET Foundation under one or more agreements.
+// The .NET Foundation licenses this file to you under the MIT License.
+// See the LICENSE file in the project root for more information.
+
+using System.Threading.Tasks;
+
+namespace System.Reactive
+{
+ public abstract class AsyncObservableBase<T> : IAsyncObservable<T>
+ {
+ public async ValueTask<IAsyncDisposable> SubscribeAsync(IAsyncObserver<T> observer)
+ {
+ if (observer == null)
+ throw new ArgumentNullException(nameof(observer));
+
+ var autoDetach = new AutoDetachAsyncObserver(observer);
+
+ var subscription = await SubscribeAsyncCore(autoDetach).ConfigureAwait(false);
+
+ await autoDetach.AssignAsync(subscription).ConfigureAwait(false);
+
+ return autoDetach;
+ }
+
+ protected abstract ValueTask<IAsyncDisposable> SubscribeAsyncCore(IAsyncObserver<T> observer);
+
+ private sealed class AutoDetachAsyncObserver : AsyncObserverBase<T>, IAsyncDisposable
+ {
+ private readonly IAsyncObserver<T> _observer;
+ private readonly object _gate = new object();
+
+ private IAsyncDisposable _subscription;
+ private ValueTask _task;
+ private bool _disposing;
+
+ public AutoDetachAsyncObserver(IAsyncObserver<T> observer)
+ {
+ _observer = observer;
+ }
+
+ public async ValueTask AssignAsync(IAsyncDisposable subscription)
+ {
+ var shouldDispose = false;
+
+ lock (_gate)
+ {
+ if (_disposing)
+ {
+ shouldDispose = true;
+ }
+ else
+ {
+ _subscription = subscription;
+ }
+ }
+
+ if (shouldDispose)
+ {
+ await subscription.DisposeAsync().ConfigureAwait(false);
+ }
+ }
+
+ public async ValueTask DisposeAsync()
+ {
+ var task = default(ValueTask);
+ var subscription = default(IAsyncDisposable);
+
+ lock (_gate)
+ {
+ //
+ // NB: The postcondition of awaiting the first DisposeAsync call to complete is that all message
+ // processing has ceased, i.e. no further On*AsyncCore calls will be made. This is achieved
+ // here by setting _disposing to true, which is checked by the On*AsyncCore calls upon
+ // entry, and by awaiting the task of any in-flight On*AsyncCore calls.
+ //
+ // Timing of the disposal of the subscription is less deterministic due to the intersection
+ // with the AssignAsync code path. However, the auto-detach observer can only be returned
+ // from the SubscribeAsync call *after* a call to AssignAsync has been made and awaited, so
+ // either AssignAsync triggers the disposal and an already disposed instance is returned, or
+ // the user calling DisposeAsync will either encounter a busy observer which will be stopped
+ // in its tracks (as described above) or it will trigger a disposal of the subscription. In
+ // both these cases the result of awaiting DisposeAsync guarantees no further message flow.
+ //
+
+ if (!_disposing)
+ {
+ _disposing = true;
+
+ task = _task;
+ subscription = _subscription;
+ }
+ }
+
+ try
+ {
+ //
+ // BUGBUG: This causes grief when an outgoing On*Async call reenters the DisposeAsync method and
+ // results in the task returned from the On*Async call to be awaited to serialize the
+ // call to subscription.DisposeAsync after it's done. We need to either detect reentrancy
+ // and queue up the call to DisposeAsync or follow an when we trigger the disposal without
+ // awaiting outstanding work (thus allowing for concurrency).
+ //
+ // if (task != null)
+ // {
+ // await task.ConfigureAwait(false);
+ // }
+ //
+ }
+ finally
+ {
+ if (subscription != null)
+ {
+ await subscription.DisposeAsync().ConfigureAwait(false);
+ }
+ }
+ }
+
+ protected override async ValueTask OnCompletedAsyncCore()
+ {
+ lock (_gate)
+ {
+ if (_disposing)
+ {
+ return;
+ }
+
+ _task = _observer.OnCompletedAsync();
+ }
+
+ try
+ {
+ await _task.ConfigureAwait(false);
+ }
+ finally
+ {
+ await FinishAsync().ConfigureAwait(false);
+ }
+ }
+
+ protected override async ValueTask OnErrorAsyncCore(Exception error)
+ {
+ lock (_gate)
+ {
+ if (_disposing)
+ {
+ return;
+ }
+
+ _task = _observer.OnErrorAsync(error);
+ }
+
+ try
+ {
+ await _task.ConfigureAwait(false);
+ }
+ finally
+ {
+ await FinishAsync().ConfigureAwait(false);
+ }
+ }
+
+ protected override async ValueTask OnNextAsyncCore(T value)
+ {
+ lock (_gate)
+ {
+ if (_disposing)
+ {
+ return;
+ }
+
+ _task = _observer.OnNextAsync(value);
+ }
+
+ try
+ {
+ await _task.ConfigureAwait(false);
+ }
+ finally
+ {
+ lock (_gate)
+ {
+ _task = default;
+ }
+ }
+ }
+
+ private async ValueTask FinishAsync()
+ {
+ var subscription = default(IAsyncDisposable);
+
+ lock (_gate)
+ {
+ if (!_disposing)
+ {
+ _disposing = true;
+
+ subscription = _subscription;
+ }
+
+ _task = default;
+ }
+
+ if (subscription != null)
+ {
+ await subscription.DisposeAsync().ConfigureAwait(false);
+ }
+ }
+ }
+ }
+}
\ No newline at end of file
--- /dev/null
+// Licensed to the .NET Foundation under one or more agreements.
+// The .NET Foundation licenses this file to you under the MIT License.
+// See the LICENSE file in the project root for more information.
+
+using System.Threading.Tasks;
+
+namespace System.Reactive
+{
+ public class AsyncObserver<T> : AsyncObserverBase<T>
+ {
+ private readonly Func<T, ValueTask> _onNextAsync;
+ private readonly Func<Exception, ValueTask> _onErrorAsync;
+ private readonly Func<ValueTask> _onCompletedAsync;
+
+ public AsyncObserver(Func<T, ValueTask> onNextAsync, Func<Exception, ValueTask> onErrorAsync, Func<ValueTask> onCompletedAsync)
+ {
+ _onNextAsync = onNextAsync ?? throw new ArgumentNullException(nameof(onNextAsync));
+ _onErrorAsync = onErrorAsync ?? throw new ArgumentNullException(nameof(onErrorAsync));
+ _onCompletedAsync = onCompletedAsync ?? throw new ArgumentNullException(nameof(onCompletedAsync));
+ }
+
+ protected override ValueTask OnCompletedAsyncCore() => _onCompletedAsync();
+
+ protected override ValueTask OnErrorAsyncCore(Exception error) => _onErrorAsync(error ?? throw new ArgumentNullException(nameof(error)));
+
+ protected override ValueTask OnNextAsyncCore(T value) => _onNextAsync(value);
+ }
+}
\ No newline at end of file
--- /dev/null
+// Licensed to the .NET Foundation under one or more agreements.
+// The .NET Foundation licenses this file to you under the MIT License.
+// See the LICENSE file in the project root for more information.
+
+using System.Threading;
+using System.Threading.Tasks;
+
+namespace System.Reactive
+{
+ public abstract class AsyncObserverBase<T> : IAsyncObserver<T>
+ {
+ private const int Idle = 0;
+ private const int Busy = 1;
+ private const int Done = 2;
+
+ private int _status = Idle;
+
+ public ValueTask OnCompletedAsync()
+ {
+ TryEnter();
+
+ try
+ {
+ return OnCompletedAsyncCore();
+ }
+ finally
+ {
+ Interlocked.Exchange(ref _status, Done);
+ }
+ }
+
+ protected abstract ValueTask OnCompletedAsyncCore();
+
+ public ValueTask OnErrorAsync(Exception error)
+ {
+ if (error == null)
+ throw new ArgumentNullException(nameof(error));
+
+ TryEnter();
+
+ try
+ {
+ return OnErrorAsyncCore(error);
+ }
+ finally
+ {
+ Interlocked.Exchange(ref _status, Done);
+ }
+ }
+
+ protected abstract ValueTask OnErrorAsyncCore(Exception error);
+
+ public ValueTask OnNextAsync(T value)
+ {
+ TryEnter();
+
+ try
+ {
+ return OnNextAsyncCore(value);
+ }
+ finally
+ {
+ Interlocked.Exchange(ref _status, Idle);
+ }
+ }
+
+ protected abstract ValueTask OnNextAsyncCore(T value);
+
+ private void TryEnter()
+ {
+ var old = Interlocked.CompareExchange(ref _status, Busy, Idle);
+
+ switch (old)
+ {
+ case Busy:
+ throw new InvalidOperationException("The observer is currently processing a notification.");
+ case Done:
+ throw new InvalidOperationException("The observer has already terminated.");
+ }
+ }
+ }
+}
\ No newline at end of file
--- /dev/null
+// Licensed to the .NET Foundation under one or more agreements.
+// The .NET Foundation licenses this file to you under the MIT License.
+// See the LICENSE file in the project root for more information.
+
+using System.Threading;
+using System.Threading.Tasks;
+
+namespace System.Reactive.Disposables
+{
+ public static class AsyncDisposable
+ {
+ public static IAsyncDisposable Nop { get; } = new NopAsyncDisposable();
+
+ public static IAsyncDisposable Create(Func<ValueTask> dispose)
+ {
+ if (dispose == null)
+ throw new ArgumentNullException(nameof(dispose));
+
+ return new AnonymousAsyncDisposable(dispose);
+ }
+
+ private sealed class AnonymousAsyncDisposable : IAsyncDisposable
+ {
+ private Func<ValueTask> _dispose;
+
+ public AnonymousAsyncDisposable(Func<ValueTask> dispose) => _dispose = dispose;
+
+ public ValueTask DisposeAsync() => Interlocked.Exchange(ref _dispose, null)?.Invoke() ?? default;
+ }
+
+ private sealed class NopAsyncDisposable : IAsyncDisposable
+ {
+ public ValueTask DisposeAsync() => default;
+ }
+ }
+}
\ No newline at end of file
--- /dev/null
+// Licensed to the .NET Foundation under one or more agreements.
+// The .NET Foundation licenses this file to you under the MIT License.
+// See the LICENSE file in the project root for more information.
+
+using System.Collections.Generic;
+using System.Linq;
+using System.Threading.Tasks;
+
+namespace System.Reactive.Subjects
+{
+ public sealed class ConcurrentSimpleAsyncSubject<T> : SimpleAsyncSubject<T>
+ {
+ protected override ValueTask OnCompletedAsyncCore(IEnumerable<IAsyncObserver<T>> observers) => new ValueTask(Task.WhenAll(observers.Select(observer => observer.OnCompletedAsync().AsTask())));
+
+ protected override ValueTask OnErrorAsyncCore(IEnumerable<IAsyncObserver<T>> observers, Exception error) => new ValueTask(Task.WhenAll(observers.Select(observer => observer.OnErrorAsync(error).AsTask())));
+
+ protected override ValueTask OnNextAsyncCore(IEnumerable<IAsyncObserver<T>> observers, T value) => new ValueTask(Task.WhenAll(observers.Select(observer => observer.OnNextAsync(value).AsTask())));
+ }
+}
\ No newline at end of file
--- /dev/null
+namespace System.Reactive.Subjects
+{
+ public interface IAsyncSubject<in TInput, out TOutput> : IAsyncObservable<TOutput>, IAsyncObserver<TInput>
+ {
+ }
+
+ public interface IAsyncSubject<T> : IAsyncSubject<T, T>
+ {
+ }
+}
\ No newline at end of file
--- /dev/null
+// Licensed to the .NET Foundation under one or more agreements.
+// The .NET Foundation licenses this file to you under the MIT License.
+// See the LICENSE file in the project root for more information.
+
+using System.Collections.Generic;
+using System.Threading.Tasks;
+
+namespace System.Reactive.Subjects
+{
+ public sealed class SequentialSimpleAsyncSubject<T> : SimpleAsyncSubject<T>
+ {
+ protected override async ValueTask OnCompletedAsyncCore(IEnumerable<IAsyncObserver<T>> observers)
+ {
+ foreach (var observer in observers)
+ {
+ await observer.OnCompletedAsync().ConfigureAwait(false);
+ }
+ }
+
+ protected override async ValueTask OnErrorAsyncCore(IEnumerable<IAsyncObserver<T>> observers, Exception error)
+ {
+ foreach (var observer in observers)
+ {
+ await observer.OnErrorAsync(error).ConfigureAwait(false);
+ }
+ }
+
+ protected override async ValueTask OnNextAsyncCore(IEnumerable<IAsyncObserver<T>> observers, T value)
+ {
+ foreach (var observer in observers)
+ {
+ await observer.OnNextAsync(value).ConfigureAwait(false);
+ }
+ }
+ }
+}
\ No newline at end of file
--- /dev/null
+// Licensed to the .NET Foundation under one or more agreements.
+// The .NET Foundation licenses this file to you under the MIT License.
+// See the LICENSE file in the project root for more information.
+
+using System.Collections.Generic;
+using System.Reactive.Disposables;
+using System.Threading.Tasks;
+
+namespace System.Reactive.Subjects
+{
+ public abstract class SimpleAsyncSubject<T> : IAsyncSubject<T>
+ {
+ private readonly object _gate = new object();
+ private readonly List<IAsyncObserver<T>> _observers = new List<IAsyncObserver<T>>();
+ private bool _done;
+ private Exception _error;
+
+ public ValueTask OnCompletedAsync()
+ {
+ IAsyncObserver<T>[] observers;
+
+ lock (_gate)
+ {
+ if (_done || _error != null)
+ {
+ return default;
+ }
+
+ _done = true;
+
+ observers = _observers.ToArray();
+ }
+
+ return OnCompletedAsyncCore(observers);
+ }
+
+ protected abstract ValueTask OnCompletedAsyncCore(IEnumerable<IAsyncObserver<T>> observers);
+
+ public ValueTask OnErrorAsync(Exception error)
+ {
+ if (error == null)
+ throw new ArgumentNullException(nameof(error));
+
+ IAsyncObserver<T>[] observers;
+
+ lock (_gate)
+ {
+ if (_done || _error != null)
+ {
+ return default;
+ }
+
+ _error = error;
+
+ observers = _observers.ToArray();
+ }
+
+ return OnErrorAsyncCore(observers, error);
+ }
+
+ protected abstract ValueTask OnErrorAsyncCore(IEnumerable<IAsyncObserver<T>> observers, Exception error);
+
+ public ValueTask OnNextAsync(T value)
+ {
+ IAsyncObserver<T>[] observers;
+
+ lock (_gate)
+ {
+ if (_done || _error != null)
+ {
+ return default;
+ }
+
+ observers = _observers.ToArray();
+ }
+
+ return OnNextAsyncCore(observers, value);
+ }
+
+ protected abstract ValueTask OnNextAsyncCore(IEnumerable<IAsyncObserver<T>> observers, T value);
+
+ public async ValueTask<IAsyncDisposable> SubscribeAsync(IAsyncObserver<T> observer)
+ {
+ if (observer == null)
+ throw new ArgumentNullException(nameof(observer));
+
+ bool done;
+ Exception error;
+
+ lock (_gate)
+ {
+ done = _done;
+ error = _error;
+
+ if (!done && error == null)
+ {
+ _observers.Add(observer);
+ }
+ }
+
+ if (done)
+ {
+ await observer.OnCompletedAsync().ConfigureAwait(false);
+
+ return AsyncDisposable.Nop;
+ }
+ else if (error != null)
+ {
+ await observer.OnErrorAsync(error).ConfigureAwait(false);
+
+ return AsyncDisposable.Nop;
+ }
+ else
+ {
+ return AsyncDisposable.Create(() =>
+ {
+ lock (_gate)
+ {
+ var i = _observers.LastIndexOf(observer);
+
+ if (i >= 0)
+ {
+ _observers.RemoveAt(i);
+ }
+ }
+
+ return default;
+ });
+ }
+ }
+ }
+}
\ No newline at end of file
--- /dev/null
+namespace Impostor.Server
+{
+ public class ClientVersionUnsupportedException : ImpostorException
+ {
+ public ClientVersionUnsupportedException(int version)
+ : base($"Version {version} is not supported by Impostor")
+ {
+ Version = version;
+ }
+
+ public int Version { get; }
+ }
+}
\ No newline at end of file
--- /dev/null
+using System;
+using System.Runtime.Serialization;
+
+namespace Impostor.Server
+{
+ public class ImpostorException : Exception
+ {
+ public ImpostorException()
+ {
+ }
+
+ protected ImpostorException(SerializationInfo info, StreamingContext context) : base(info, context)
+ {
+ }
+
+ public ImpostorException(string? message) : base(message)
+ {
+ }
+
+ public ImpostorException(string? message, Exception? innerException) : base(message, innerException)
+ {
+ }
+ }
+}
\ No newline at end of file
--- /dev/null
+using System;
+using Impostor.Shared.Innersloth;
+
+namespace Impostor.Server
+{
+ public readonly struct GameCode : IEquatable<GameCode>
+ {
+ public GameCode(int value)
+ {
+ Value = value;
+ Code = GameCodeParser.IntToGameName(value);
+ }
+
+ public GameCode(string code)
+ {
+ Value = GameCodeParser.GameNameToInt(code);
+ Code = code;
+ }
+
+ public string Code { get; }
+
+ public int Value { get; }
+
+ public static implicit operator string(GameCode code) => code.Code;
+
+ public static implicit operator int(GameCode code) => code.Value;
+
+ public static implicit operator GameCode(string code) => From(code);
+
+ public static implicit operator GameCode(int value) => From(value);
+
+ public bool Equals(GameCode other)
+ {
+ return Code == other.Code && Value == other.Value;
+ }
+
+ public override bool Equals(object? obj)
+ {
+ return obj is GameCode other && Equals(other);
+ }
+
+ public override int GetHashCode()
+ {
+ return HashCode.Combine(Code, Value);
+ }
+
+ public static bool operator ==(GameCode left, GameCode right)
+ {
+ return left.Equals(right);
+ }
+
+ public static bool operator !=(GameCode left, GameCode right)
+ {
+ return !left.Equals(right);
+ }
+
+ public override string ToString()
+ {
+ return Code;
+ }
+
+ public static GameCode From(int value) => new GameCode(value);
+
+ public static GameCode From(string value) => new GameCode(value);
+
+ public static GameCode Create()
+ {
+ return new GameCode(GameCodeParser.GenerateCode(6));
+ }
+ }
+}
\ No newline at end of file
--- /dev/null
+using System.Collections.Generic;
+using System.Diagnostics.CodeAnalysis;
+using Impostor.Server.Net;
+
+namespace Impostor.Server
+{
+ public interface IGame
+ {
+ GameCode Code { get; }
+
+ IEnumerable<IClientPlayer> Players { get; }
+
+ IClientPlayer Host { get; }
+
+ bool IsPublic { get; }
+ IDictionary<object,object> Items { get; }
+
+ IGameMessageWriter CreateMessage(MessageType type);
+
+ bool TryGetPlayer(int id, [NotNullWhen(true)] out IClientPlayer player);
+ }
+}
\ No newline at end of file
--- /dev/null
+<Project Sdk="Microsoft.NET.Sdk">
+
+ <PropertyGroup>
+ <TargetFramework>net5.0</TargetFramework>
+ <RootNamespace>Impostor.Server</RootNamespace>
+ <Nullable>enable</Nullable>
+ </PropertyGroup>
+
+ <ItemGroup>
+ <ProjectReference Include="..\Imposter.Reactive\Imposter.Reactive.csproj" />
+ <ProjectReference Include="..\Impostor.Shared\Impostor.Shared.csproj" />
+ </ItemGroup>
+
+ <ItemGroup>
+ <PackageReference Include="Microsoft.Extensions.Logging.Abstractions" Version="5.0.0-rc.1.20451.14" />
+ <PackageReference Include="Microsoft.Extensions.DependencyInjection.Abstractions" Version="5.0.0-rc.1.20451.14" />
+ <PackageReference Include="System.Interactive.Async" Version="4.1.1" />
+ </ItemGroup>
+
+</Project>
--- /dev/null
+<wpf:ResourceDictionary xml:space="preserve" xmlns:x="http://schemas.microsoft.com/winfx/2006/xaml" xmlns:s="clr-namespace:System;assembly=mscorlib" xmlns:ss="urn:shemas-jetbrains-com:settings-storage-xaml" xmlns:wpf="http://schemas.microsoft.com/winfx/2006/xaml/presentation">
+ <s:Boolean x:Key="/Default/CodeInspection/NamespaceProvider/NamespaceFoldersToSkip/=exceptions/@EntryIndexedValue">True</s:Boolean>
+ <s:Boolean x:Key="/Default/CodeInspection/NamespaceProvider/NamespaceFoldersToSkip/=net_005Cextensions/@EntryIndexedValue">True</s:Boolean></wpf:ResourceDictionary>
\ No newline at end of file
--- /dev/null
+using System;
+
+namespace Impostor.Shared.Innersloth.Data
+{
+ [Flags]
+ public enum LimboStates
+ {
+ PreSpawn = 1,
+ NotLimbo = 2,
+ WaitingForHost = 4,
+ All = PreSpawn | NotLimbo | WaitingForHost
+ }
+}
\ No newline at end of file
--- /dev/null
+using System;
+using System.Threading.Tasks;
+using Impostor.Shared.Innersloth.Data;
+
+namespace Impostor.Server.Net
+{
+ public static class GameMessageWriterExtensions
+ {
+ public static ValueTask SendToAllExceptAsync(this IGameMessageWriter writer, LimboStates states, int? id)
+ {
+ return id.HasValue
+ ? writer.SendToAllExceptAsync(states, id.Value)
+ : writer.SendToAllAsync(states);
+ }
+
+ public static ValueTask SendToAllExceptAsync(this IGameMessageWriter writer, LimboStates states, IClient client)
+ {
+ if (client == null) throw new ArgumentNullException(nameof(client));
+
+ return writer.SendToAllExceptAsync(states, client.Id);
+ }
+
+ public static ValueTask SendToAsync(this IGameMessageWriter writer, IClient client)
+ {
+ if (client == null) throw new ArgumentNullException(nameof(client));
+
+ return writer.SendToAsync(client.Id);
+ }
+ }
+}
\ No newline at end of file
--- /dev/null
+using System.Threading.Tasks;
+
+namespace Impostor.Server.Net.Factories
+{
+ public interface IClientFactory
+ {
+ /// <summary>
+ /// Get the next ID for <see cref="IClient"/>.
+ /// </summary>
+ int NextId();
+
+ /// <summary>
+ /// Creates a client for the Hazel <see cref="connection"/>.
+ /// </summary>
+ /// <param name="connection">Hazel connection.</param>
+ /// <param name="name"></param>
+ /// <param name="clientVersion"></param>
+ ValueTask<IClient> CreateAsync(IConnection connection, string name, int clientVersion);
+ }
+}
\ No newline at end of file
--- /dev/null
+using System.Collections.Generic;
+
+namespace Impostor.Server.Net
+{
+ public interface IClient
+ {
+ int Id { get; }
+
+ string Name { get; }
+
+ IConnection Connection { get; }
+
+ IDictionary<object,object> Items { get; }
+ }
+}
\ No newline at end of file
--- /dev/null
+using Impostor.Shared.Innersloth.Data;
+
+namespace Impostor.Server.Net
+{
+ public interface IClientPlayer
+ {
+ IClient Client { get; }
+
+ IGame Game { get; }
+
+ LimboStates Limbo { get; }
+ }
+}
\ No newline at end of file
--- /dev/null
+using System;
+using System.Net;
+
+namespace Impostor.Server.Net
+{
+ public interface IConnection
+ {
+ IAsyncObservable<IMessage> MessageReceived { get; }
+
+ IPEndPoint EndPoint { get; }
+
+ bool IsConnected { get; }
+
+ IConnectionMessageWriter CreateMessage(MessageType type);
+ }
+}
\ No newline at end of file
--- /dev/null
+using System.Threading.Tasks;
+
+namespace Impostor.Server.Net.Manager
+{
+ public interface IClientManager
+ {
+ ValueTask RegisterConnectionAsync(IConnection connection, string name, int clientVersion);
+
+ void Register(IClient client);
+
+ void Remove(IClient client);
+ }
+}
\ No newline at end of file
--- /dev/null
+using System.Net;
+using System.Threading.Tasks;
+
+namespace Impostor.Server.Net.Manager
+{
+ public interface IMatchmaker
+ {
+ ValueTask StartAsync(IPEndPoint ipEndPoint);
+
+ ValueTask StopAsync();
+
+ IGameMessageWriter CreateGameMessageWriter(IGame game, MessageType messageType);
+ }
+}
\ No newline at end of file
--- /dev/null
+using System.Threading.Tasks;
+
+namespace Impostor.Server.Net
+{
+ public interface IConnectionMessageWriter : IMessageWriter
+ {
+ ValueTask SendAsync();
+ }
+}
\ No newline at end of file
--- /dev/null
+using System.Threading.Tasks;
+using Impostor.Shared.Innersloth.Data;
+
+namespace Impostor.Server.Net
+{
+ public interface IGameMessageWriter : IMessageWriter
+ {
+ /// <summary>
+ /// Send the message to all players.
+ /// </summary>
+ /// <param name="states"></param>
+ ValueTask SendToAllAsync(LimboStates states);
+
+ /// <summary>
+ /// Send the message to all players except one.
+ /// </summary>
+ /// <param name="states"></param>
+ /// <param name="senderId">The player to exclude from sending the message.</param>
+ ValueTask SendToAllExceptAsync(LimboStates states, int senderId);
+
+ /// <summary>
+ /// Send a message to a specific player.
+ /// </summary>
+ /// <param name="id"></param>
+ ValueTask SendToAsync(int id);
+ }
+}
\ No newline at end of file
--- /dev/null
+namespace Impostor.Server.Net
+{
+ public interface IMessage
+ {
+ MessageType Type { get; }
+
+ IMessageReader CreateReader();
+ }
+}
\ No newline at end of file
--- /dev/null
+using System;
+
+namespace Impostor.Server.Net
+{
+ public interface IMessageReader
+ {
+ int Position { get; }
+
+ ReadOnlyMemory<byte> Buffer { get; }
+
+ byte Tag { get; }
+
+ int Length { get; }
+
+ bool ReadBoolean();
+
+ sbyte ReadSByte();
+
+ byte ReadByte();
+
+ ushort ReadUInt16();
+
+ short ReadInt16();
+
+ uint ReadUInt32();
+
+ int ReadInt32();
+
+ float ReadSingle();
+
+ string ReadString();
+
+ ReadOnlyMemory<byte> ReadBytesAndSize();
+
+ ReadOnlyMemory<byte> ReadBytes(int length);
+
+ int ReadPackedInt32();
+
+ uint ReadPackedUInt32();
+
+ void CopyTo(IMessageWriter writer);
+ }
+}
\ No newline at end of file
--- /dev/null
+using System;
+using System.Net;
+
+namespace Impostor.Server.Net
+{
+ public interface IMessageWriter : IDisposable
+ {
+ void Write(bool value);
+
+ void Write(sbyte value);
+
+ void Write(byte value);
+
+ void Write(short value);
+
+ void Write(ushort value);
+
+ void Write(uint value);
+
+ void Write(int value);
+
+ void Write(float value);
+
+ void Write(string value);
+
+ void Write(IPAddress ipAddress);
+
+ void WritePacked(int value);
+
+ void Write(ReadOnlyMemory<byte> data);
+
+ void StartMessage(byte typeFlag);
+
+ void Write(GameCode code);
+
+ void EndMessage();
+
+ void Clear(MessageType type);
+ }
+}
\ No newline at end of file
--- /dev/null
+namespace Impostor.Server.Net
+{
+ public enum MessageType
+ {
+ Unreliable,
+ Reliable
+ }
+}
\ No newline at end of file
--- /dev/null
+using Impostor.Server.Net.Manager;
+using Microsoft.Extensions.DependencyInjection;
+
+namespace Impostor.Server.Hazel
+{
+ public static class ServiceExtensions
+ {
+ public static IServiceCollection UseHazelMatchmaking(this IServiceCollection services)
+ {
+ services.AddSingleton<IMatchmaker, HazelMatchmaker>();
+ return services;
+ }
+ }
+}
\ No newline at end of file
--- /dev/null
+using System;
+using System.Net;
+using System.Reactive.Subjects;
+using System.Threading.Tasks;
+using Hazel;
+using Impostor.Server.Net;
+using Microsoft.Extensions.Logging;
+
+namespace Impostor.Server.Hazel
+{
+ internal class HazelConnection : IConnection
+ {
+ private readonly ILogger<HazelConnection> _logger;
+ private readonly ConcurrentSimpleAsyncSubject<IMessage> _messageReceived = new ConcurrentSimpleAsyncSubject<IMessage>();
+
+ public HazelConnection(Connection innerConnection, ILogger<HazelConnection> logger)
+ {
+ _logger = logger;
+ InnerConnection = innerConnection;
+ innerConnection.DataReceived += ConnectionOnDataReceived;
+ innerConnection.Disconnected += ConnectionOnDisconnected;
+ }
+
+ public Connection InnerConnection { get; }
+
+ public IAsyncObservable<IMessage> MessageReceived => _messageReceived;
+
+ public IPEndPoint EndPoint => InnerConnection.EndPoint;
+
+ public bool IsConnected => InnerConnection.State == ConnectionState.Connected;
+
+ private void ConnectionOnDisconnected(object sender, DisconnectedEventArgs e)
+ {
+ Task.Run(_messageReceived.OnCompletedAsync);
+ }
+
+ private void ConnectionOnDataReceived(DataReceivedEventArgs e)
+ {
+ Task.Run(() => HandleData(e));
+ }
+
+ private async Task HandleData(DataReceivedEventArgs e)
+ {
+ try
+ {
+ while (true)
+ {
+ if (e.Message.Position >= e.Message.Length)
+ {
+ break;
+ }
+
+ var reader = e.Message.ReadMessage();
+ var type = e.SendOption switch
+ {
+ SendOption.None => MessageType.Unreliable,
+ SendOption.Reliable => MessageType.Reliable,
+ _ => throw new NotSupportedException()
+ };
+
+ using var message = new HazelMessage(reader, type);
+
+ await _messageReceived.OnNextAsync(message);
+ }
+ }
+ catch (Exception ex)
+ {
+ _logger.LogError(ex, "Exception caught in client data handler.");
+ }
+ finally
+ {
+ e.Message.Recycle();
+ }
+ }
+
+ public IConnectionMessageWriter CreateMessage(MessageType type)
+ {
+ return new HazelConnectionMessageWriter(type, InnerConnection);
+ }
+ }
+}
\ No newline at end of file
--- /dev/null
+using System;
+using System.Net;
+using System.Net.Sockets;
+using System.Threading.Tasks;
+using Hazel;
+using Hazel.Udp;
+using Impostor.Server.Net;
+using Impostor.Server.Net.Manager;
+using Microsoft.Extensions.Logging;
+
+namespace Impostor.Server.Hazel
+{
+ internal class HazelMatchmaker : IMatchmaker
+ {
+ private readonly IClientManager _clientManager;
+ private readonly ILogger<HazelMatchmaker> _logger;
+ private readonly ILogger<HazelConnection> _connectionLogger;
+ private UdpConnectionListener _connection;
+
+ public HazelMatchmaker(
+ ILogger<HazelMatchmaker> logger,
+ IClientManager clientManager,
+ ILogger<HazelConnection> connectionLogger)
+ {
+ _logger = logger;
+ _clientManager = clientManager;
+ _connectionLogger = connectionLogger;
+ }
+
+ public ValueTask StartAsync(IPEndPoint ipEndPoint)
+ {
+ var mode = ipEndPoint.AddressFamily switch
+ {
+ AddressFamily.InterNetwork => IPMode.IPv4,
+ AddressFamily.InterNetworkV6 => IPMode.IPv6,
+ _ => throw new InvalidOperationException()
+ };
+
+ _connection = new UdpConnectionListener(ipEndPoint, mode, s =>
+ {
+ _logger.LogWarning("Log from Hazel: {0}", s);
+ });
+
+ _connection.NewConnection += OnNewConnection;
+
+ _connection.Start();
+
+ return default;
+ }
+
+ public ValueTask StopAsync()
+ {
+ _connection.Dispose();
+
+ return default;
+ }
+
+ private void OnNewConnection(NewConnectionEventArgs e)
+ {
+ Task.Run(() => HandleNewConnection(e));
+ }
+
+ private async Task HandleNewConnection(NewConnectionEventArgs e)
+ {
+ int clientVersion;
+ string name;
+ try
+ {
+ // Handshake.
+ clientVersion = e.HandshakeData.ReadInt32();
+ name = e.HandshakeData.ReadString();
+
+ e.HandshakeData.Recycle();
+ }
+ catch (Exception ex)
+ {
+ _logger.LogTrace(ex, "Error in new connection.");
+ return;
+ }
+
+ var connection = new HazelConnection(e.Connection, _connectionLogger);
+
+ // Register client
+ await _clientManager.RegisterConnectionAsync(connection, name, clientVersion);
+ }
+
+ public IGameMessageWriter CreateGameMessageWriter(IGame game, MessageType messageType)
+ {
+ return new HazelGameMessageWriter(messageType, game);
+ }
+ }
+}
\ No newline at end of file
--- /dev/null
+<Project Sdk="Microsoft.NET.Sdk">
+
+ <PropertyGroup>
+ <TargetFramework>net5.0</TargetFramework>
+ <AllowUnsafeBlocks>true</AllowUnsafeBlocks>
+ </PropertyGroup>
+
+ <ItemGroup>
+ <ProjectReference Include="..\..\submodules\Hazel-Networking\Hazel\Hazel.csproj" />
+ <ProjectReference Include="..\Impostor.Server.Api\Impostor.Server.Api.csproj" />
+ </ItemGroup>
+
+ <ItemGroup>
+ <PackageReference Include="Microsoft.Extensions.DependencyInjection.Abstractions" Version="5.0.0-rc.1.20451.14" />
+ </ItemGroup>
+
+</Project>
--- /dev/null
+<wpf:ResourceDictionary xml:space="preserve" xmlns:x="http://schemas.microsoft.com/winfx/2006/xaml" xmlns:s="clr-namespace:System;assembly=mscorlib" xmlns:ss="urn:shemas-jetbrains-com:settings-storage-xaml" xmlns:wpf="http://schemas.microsoft.com/winfx/2006/xaml/presentation">
+ <s:Boolean x:Key="/Default/CodeInspection/NamespaceProvider/NamespaceFoldersToSkip/=extensions/@EntryIndexedValue">True</s:Boolean></wpf:ResourceDictionary>
\ No newline at end of file
--- /dev/null
+using System;
+using System.Runtime.CompilerServices;
+using System.Text;
+using Impostor.Server.Net;
+
+namespace Impostor.Server.Hazel
+{
+ public class BufferMessageReader : IMessageReader
+ {
+ // TODO: Remove _offset, we can slice the buffer.
+ private readonly int _offset;
+ private int _position;
+ private int readHead;
+
+ public ReadOnlyMemory<byte> Buffer { get; }
+
+ public byte Tag { get; }
+
+ public int Length { get; }
+
+ public int Position
+ {
+ get { return _position; }
+ set
+ {
+ _position = value;
+ readHead = value + _offset;
+ }
+ }
+
+ public BufferMessageReader(byte tag, ReadOnlyMemory<byte> buffer, int offset, int length)
+ {
+ Tag = tag;
+ Buffer = buffer;
+ Length = length;
+ _offset = offset;
+ readHead = offset;
+ }
+
+ public bool ReadBoolean()
+ {
+ byte val = FastByte();
+ return val != 0;
+ }
+
+ public sbyte ReadSByte()
+ {
+ return (sbyte)FastByte();
+ }
+
+ public byte ReadByte()
+ {
+ return FastByte();
+ }
+
+ public ushort ReadUInt16()
+ {
+ // TODO: Refactor to System.Buffers.Binary.BinaryPrimitives
+
+ ushort output =
+ (ushort)(FastByte()
+ | FastByte() << 8);
+ return output;
+ }
+
+ public short ReadInt16()
+ {
+ // TODO: Refactor to System.Buffers.Binary.BinaryPrimitives
+
+ short output =
+ (short)(FastByte() | FastByte() << 8);
+ return output;
+ }
+
+ public uint ReadUInt32()
+ {
+ // TODO: Refactor to System.Buffers.Binary.BinaryPrimitives
+
+ uint output = FastByte()
+ | (uint)FastByte() << 8
+ | (uint)FastByte() << 16
+ | (uint)FastByte() << 24;
+
+ return output;
+ }
+
+ public int ReadInt32()
+ {
+ // TODO: Refactor to System.Buffers.Binary.BinaryPrimitives
+
+ int output = FastByte()
+ | FastByte() << 8
+ | FastByte() << 16
+ | FastByte() << 24;
+
+ return output;
+ }
+
+ public unsafe float ReadSingle()
+ {
+ // TODO: Refactor to System.Buffers.Binary.BinaryPrimitives
+
+ float output = 0;
+ fixed (byte* bufPtr = &Buffer.Span[readHead])
+ {
+ byte* outPtr = (byte*)&output;
+
+ *outPtr = *bufPtr;
+ *(outPtr + 1) = *(bufPtr + 1);
+ *(outPtr + 2) = *(bufPtr + 2);
+ *(outPtr + 3) = *(bufPtr + 3);
+ }
+
+ Position += 4;
+ return output;
+ }
+
+ public string ReadString()
+ {
+ var len = ReadPackedInt32();
+ var output = Encoding.UTF8.GetString(Buffer.Span.Slice(readHead, len));
+ Position += len;
+ return output;
+ }
+
+ public ReadOnlyMemory<byte> ReadBytesAndSize()
+ {
+ var len = ReadPackedInt32();
+ return ReadBytes(len);
+ }
+
+ public ReadOnlyMemory<byte> ReadBytes(int length)
+ {
+ var output = Buffer.Slice(readHead, length);
+ Position += length;
+ return output;
+ }
+
+ public int ReadPackedInt32()
+ {
+ return (int)ReadPackedUInt32();
+ }
+
+ public uint ReadPackedUInt32()
+ {
+ bool readMore = true;
+ int shift = 0;
+ uint output = 0;
+
+ while (readMore)
+ {
+ byte b = ReadByte();
+ if (b >= 0x80)
+ {
+ readMore = true;
+ b ^= 0x80;
+ }
+ else
+ {
+ readMore = false;
+ }
+
+ output |= (uint)(b << shift);
+ shift += 7;
+ }
+
+ return output;
+ }
+
+ public void CopyTo(IMessageWriter writer)
+ {
+ int offset, length;
+ if (Tag == byte.MaxValue)
+ {
+ offset = _offset;
+ length = Length;
+ }
+ else
+ {
+ offset = _offset - 3;
+ length = Length + 3;
+ }
+
+ writer.Write(Buffer.Slice(offset, length));
+ }
+
+ [MethodImpl(MethodImplOptions.AggressiveInlining)]
+ private byte FastByte()
+ {
+ _position++;
+ return Buffer.Span[readHead++];
+ }
+ }
+}
\ No newline at end of file
--- /dev/null
+using System.Threading.Tasks;
+using Hazel;
+using Impostor.Server.Net;
+
+namespace Impostor.Server.Hazel
+{
+ internal class HazelConnectionMessageWriter : HazelMessageWriter, IConnectionMessageWriter
+ {
+ private readonly Connection _connection;
+
+ public HazelConnectionMessageWriter(MessageType type, Connection connection)
+ : base(type)
+ {
+ _connection = connection;
+ }
+
+ public ValueTask SendAsync()
+ {
+ _connection.Send(Writer);
+ return default;
+ }
+ }
+}
\ No newline at end of file
--- /dev/null
+using System;
+using System.Collections.Generic;
+using System.Linq;
+using System.Threading.Tasks;
+using Hazel;
+using Impostor.Server.Net;
+using Impostor.Shared.Innersloth.Data;
+
+namespace Impostor.Server.Hazel
+{
+ internal class HazelGameMessageWriter : HazelMessageWriter, IGameMessageWriter
+ {
+ private readonly IGame _game;
+
+ public HazelGameMessageWriter(MessageType type, IGame game)
+ : base(type)
+ {
+ _game = game;
+ }
+
+ private IEnumerable<Connection> GetConnections(Func<IClientPlayer, bool> filter)
+ {
+ return _game.Players
+ .Where(filter)
+ .Select(p => p.Client.Connection)
+ .OfType<HazelConnection>()
+ .Select(c => c.InnerConnection);
+ }
+
+ public ValueTask SendToAllAsync(LimboStates states)
+ {
+ foreach (var connection in GetConnections(x => x.Limbo.HasFlag(states)))
+ {
+ connection.Send(Writer);
+ }
+
+ return default;
+ }
+
+ public ValueTask SendToAllExceptAsync(LimboStates states, int senderId)
+ {
+ foreach (var connection in GetConnections(x =>
+ x.Limbo.HasFlag(states) &&
+ x.Client.Id != senderId))
+ {
+ connection.Send(Writer);
+ }
+ return default;
+ }
+
+ public ValueTask SendToAsync(int id)
+ {
+ if (_game.TryGetPlayer(id, out var player))
+ {
+ ((HazelConnection)player.Client.Connection).InnerConnection.Send(Writer);
+ }
+
+ return default;
+ }
+ }
+}
\ No newline at end of file
--- /dev/null
+using System;
+using Hazel;
+using Impostor.Server.Net;
+
+namespace Impostor.Server.Hazel
+{
+ internal class HazelMessage : IMessage, IDisposable
+ {
+ private bool _isDisposed;
+ private readonly MessageReader _reader;
+
+ public HazelMessage(MessageReader reader, MessageType type)
+ {
+ _reader = reader;
+ Type = type;
+ }
+
+ public MessageType Type { get; }
+
+ public IMessageReader CreateReader()
+ {
+ if (_isDisposed)
+ {
+ throw new ObjectDisposedException(nameof(_reader));
+ }
+
+ return new BufferMessageReader(_reader.Tag, _reader.Buffer, _reader.Offset, _reader.Length);
+ }
+
+ private void Dispose(bool disposing)
+ {
+ if (disposing)
+ {
+ _isDisposed = true;
+ }
+ }
+
+ public void Dispose()
+ {
+ Dispose(true);
+ GC.SuppressFinalize(this);
+ }
+
+ ~HazelMessage()
+ {
+ Dispose(false);
+ }
+ }
+}
\ No newline at end of file
--- /dev/null
+using System;
+using System.Net;
+using Hazel;
+using Impostor.Server.Net;
+
+namespace Impostor.Server.Hazel
+{
+ internal abstract class HazelMessageWriter : IMessageWriter
+ {
+ protected readonly MessageWriter Writer;
+
+ protected HazelMessageWriter(MessageType type)
+ {
+ Writer = MessageWriter.Get(ToSendOption(type));
+ }
+
+ private static SendOption ToSendOption(MessageType type)
+ {
+ return type switch
+ {
+ MessageType.Unreliable => SendOption.None,
+ MessageType.Reliable => SendOption.Reliable,
+ _ => throw new NotSupportedException($"Message type {type} is not supported")
+ };
+ }
+
+ protected virtual void Dispose(bool disposing)
+ {
+ if (disposing)
+ {
+ Writer.Recycle();
+ }
+ }
+
+ public void Dispose()
+ {
+ Dispose(true);
+ GC.SuppressFinalize(this);
+ }
+
+ public void Write(bool value)
+ {
+ Writer.Write(value);
+ }
+
+ public void Write(sbyte value)
+ {
+ Writer.Write(value);
+ }
+
+ public void Write(byte value)
+ {
+ Writer.Write(value);
+ }
+
+ public void Write(short value)
+ {
+ Writer.Write(value);
+ }
+
+ public void Write(ushort value)
+ {
+ Writer.Write(value);
+ }
+
+ public void Write(uint value)
+ {
+ Writer.Write(value);
+ }
+
+ public void Write(int value)
+ {
+ Writer.Write(value);
+ }
+
+ public void Write(float value)
+ {
+ Writer.Write(value);
+ }
+
+ public void Write(string value)
+ {
+ Writer.Write(value);
+ }
+
+ public void Write(IPAddress ipAddress)
+ {
+ Writer.Write(ipAddress.GetAddressBytes());
+ }
+
+ public void WritePacked(int value)
+ {
+ Writer.WritePacked(value);
+ }
+
+ public void Write(ReadOnlyMemory<byte> data)
+ {
+ Writer.Write(data.ToArray()); // TODO: Fix memory allocation.
+ }
+
+ public void StartMessage(byte typeFlag)
+ {
+ Writer.StartMessage(typeFlag);
+ }
+
+ public void Write(GameCode code)
+ {
+ Write(code.Value);
+ }
+
+ public void EndMessage()
+ {
+ Writer.EndMessage();
+ }
+
+ public void Clear(MessageType type)
+ {
+ Writer.Clear(ToSendOption(type));
+ }
+ }
+}
\ No newline at end of file
</PropertyGroup>
<ItemGroup>
- <ProjectReference Include="..\..\submodules\Hazel-Networking\Hazel\Hazel.csproj" />
+ <ProjectReference Include="..\Impostor.Server.Api\Impostor.Server.Api.csproj" />
+ <ProjectReference Include="..\Impostor.Server.Hazel\Impostor.Server.Hazel.csproj" />
<ProjectReference Include="..\Impostor.Shared\Impostor.Shared.csproj" />
</ItemGroup>
</Content>
</ItemGroup>
+ <ItemGroup>
+ <Folder Include="Net\Hazel" />
+ </ItemGroup>
+
</Project>
using System;
-using Hazel;
-using Impostor.Server.Data;
+using System.Threading.Tasks;
using Impostor.Server.Net.Manager;
using Impostor.Server.Net.Messages;
using Impostor.Server.Net.State;
namespace Impostor.Server.Net
{
- internal class Client
+ internal class Client : ClientBase
{
private static readonly ILogger Logger = Log.ForContext<Client>();
- private readonly ClientManager _clientManager;
+ private readonly IClientManager _clientManager;
private readonly GameManager _gameManager;
- public Client(ClientManager clientManager, GameManager gameManager, int id, string name, Connection connection)
+ public Client(IClientManager clientManager, GameManager gameManager, int id, string name, IConnection connection)
+ : base(id, name, connection)
{
_clientManager = clientManager;
_gameManager = gameManager;
- Id = id;
- Name = name;
- Connection = connection;
- Connection.DataReceived += OnDataReceived;
- Connection.Disconnected += OnDisconnected;
Player = new ClientPlayer(this, _gameManager);
}
- public int Id { get; }
- public string Name { get; }
- public Connection Connection { get; }
public ClientPlayer Player { get; }
-
- public void Send(MessageWriter writer)
- {
- Connection.Send(writer);
- }
- private bool IsPacketAllowed(MessageReader message, bool hostOnly)
+ private bool IsPacketAllowed(IMessageReader message, bool hostOnly)
{
var game = Player.Game;
if (game == null)
return true;
}
- private void OnDataReceived(DataReceivedEventArgs e)
- {
- try
- {
- while (true)
- {
- if (e.Message.Position >= e.Message.Length)
- {
- break;
- }
-
- OnMessageReceived(e.Message.ReadMessage(), e.SendOption);
- }
- }
- catch (Exception ex)
- {
- Logger.Error(ex, "Exception caught in client data handler.");
- Player.SendDisconnectReason(DisconnectReason.Custom, DisconnectMessages.Error);
- }
- finally
- {
- e.Message.Recycle();
- }
- }
-
- private void OnMessageReceived(MessageReader message, SendOption sendOption)
+
+ protected override async ValueTask OnMessageReceived(IMessage message)
{
- var flag = message.Tag;
+ var reader = message.CreateReader();
+
+ var flag = reader.Tag;
Logger.Verbose("[{0}] Server got {1}.", Id, flag);
case MessageFlags.HostGame:
{
// Read game settings.
- var gameInfo = Message00HostGame.Deserialize(message);
+ var gameInfo = Message00HostGame.Deserialize(reader);
// Create game.
var game = _gameManager.Create(gameInfo);
if (game == null)
{
- Player.SendDisconnectReason(DisconnectReason.ServerFull);
+ await Player.SendDisconnectReason(DisconnectReason.ServerFull);
return;
}
// Code in the packet below will be used in JoinGame.
- using (var writer = MessageWriter.Get(SendOption.Reliable))
+ using (var writer = Connection.CreateMessage(MessageType.Reliable))
{
Message00HostGame.Serialize(writer, game.Code);
- Connection.Send(writer);
+ await writer.SendAsync();
}
break;
}
case MessageFlags.JoinGame:
{
- Message01JoinGame.Deserialize(message,
+ Message01JoinGame.Deserialize(reader,
out var gameCode,
out var unknown);
var game = _gameManager.Find(gameCode);
if (game == null)
{
- Player.SendDisconnectReason(DisconnectReason.GameMissing);
+ await Player.SendDisconnectReason(DisconnectReason.GameMissing);
return;
}
- game.HandleJoinGame(Player);
+ await game.HandleJoinGame(Player);
break;
}
case MessageFlags.StartGame:
{
- if (!IsPacketAllowed(message, true))
+ if (!IsPacketAllowed(reader, true))
{
return;
}
- Player.Game.HandleStartGame(message);
+ await Player.Game.HandleStartGame(reader);
break;
}
case MessageFlags.RemovePlayer:
{
- if (!IsPacketAllowed(message, true))
+ if (!IsPacketAllowed(reader, true))
{
return;
}
- Message04RemovePlayer.Deserialize(message,
+ Message04RemovePlayer.Deserialize(reader,
out var playerId,
out var reason);
- Player.Game.HandleRemovePlayer(playerId, (DisconnectReason) reason);
+ await Player.Game.HandleRemovePlayer(playerId, (DisconnectReason) reason);
break;
}
case MessageFlags.GameData:
case MessageFlags.GameDataTo:
{
- if (!IsPacketAllowed(message, false))
+ if (!IsPacketAllowed(reader, false))
{
return;
}
// Broadcast packet to all other players.
- using (var writer = MessageWriter.Get(sendOption))
+ using var writer = Player.Game.CreateMessage(message.Type);
+
+ if (flag == MessageFlags.GameDataTo)
{
- if (flag == MessageFlags.GameDataTo)
- {
- var target = message.ReadPackedInt32();
- writer.CopyFrom(message);
- Player.Game.SendTo(writer, target);
- }
- else
- {
- writer.CopyFrom(message);
- Player.Game.SendToAllExcept(writer, Player.Client.Id);
- }
+ var target = reader.ReadPackedInt32();
+ reader.CopyTo(writer);
+ await writer.SendToAsync(target);
}
+ else
+ {
+ reader.CopyTo(writer);
+ await writer.SendToAllExceptAsync(LimboStates.NotLimbo, Player.Client.Id);
+ }
+
break;
}
case MessageFlags.EndGame:
{
- if (!IsPacketAllowed(message, true))
+ if (!IsPacketAllowed(reader, true))
{
return;
}
- Player.Game.HandleEndGame(message);
+ await Player.Game.HandleEndGame(reader);
break;
}
case MessageFlags.AlterGame:
{
- if (!IsPacketAllowed(message, true))
+ if (!IsPacketAllowed(reader, true))
{
return;
}
- Message10AlterGame.Deserialize(message,
+ Message10AlterGame.Deserialize(reader,
out var gameTag,
out var value);
return;
}
- Player.Game.HandleAlterGame(message, Player, value);
+ await Player.Game.HandleAlterGame(reader, Player, value);
break;
}
case MessageFlags.KickPlayer:
{
- if (!IsPacketAllowed(message, true))
+ if (!IsPacketAllowed(reader, true))
{
return;
}
- Message11KickPlayer.Deserialize(message,
+ Message11KickPlayer.Deserialize(reader,
out var playerId,
out var isBan);
- Player.Game.HandleKickPlayer(playerId, isBan);
+ await Player.Game.HandleKickPlayer(playerId, isBan);
break;
}
case MessageFlags.GetGameListV2:
{
- Message16GetGameListV2.Deserialize(message, out var options);
- Player.OnRequestGameList(options);
+ Message16GetGameListV2.Deserialize(reader, out var options);
+ await Player.OnRequestGameList(options);
break;
}
if (flag != MessageFlags.GameData &&
flag != MessageFlags.GameDataTo &&
flag != MessageFlags.EndGame &&
- message.Position < message.Length)
+ reader.Position < reader.Length)
{
Logger.Warning("Server did not consume all bytes from {0} ({1} < {2}).",
flag,
- message.Position,
- message.Length);
+ reader.Position,
+ reader.Length);
}
#endif
}
- private void OnDisconnected(object sender, DisconnectedEventArgs e)
+ protected override async ValueTask OnDisconnected()
{
try
{
if (Player.Game != null)
{
- Player.Game.HandleRemovePlayer(Id, DisconnectReason.ExitGame);
+ await Player.Game.HandleRemovePlayer(Id, DisconnectReason.ExitGame);
}
}
catch (Exception ex)
--- /dev/null
+using System;
+using System.Collections.Concurrent;
+using System.Collections.Generic;
+using System.Threading.Tasks;
+
+namespace Impostor.Server.Net
+{
+ public abstract class ClientBase : IClient
+ {
+ protected ClientBase(int id, string name, IConnection connection)
+ {
+ Id = id;
+ Name = name;
+ Connection = connection;
+ Items = new ConcurrentDictionary<object, object>();
+ }
+
+ public int Id { get; }
+
+ public string Name { get; }
+
+ public IConnection Connection { get; }
+
+ public IDictionary<object, object> Items { get; }
+
+ public virtual async ValueTask InitializeAsync()
+ {
+ await Connection.MessageReceived.SubscribeAsync(OnMessageReceived, OnDisconnected);
+ }
+
+ protected abstract ValueTask OnMessageReceived(IMessage message);
+
+ protected abstract ValueTask OnDisconnected();
+ }
+}
\ No newline at end of file
--- /dev/null
+using System;
+using System.Threading;
+using System.Threading.Tasks;
+using Impostor.Server.Net.Messages;
+using Impostor.Shared.Innersloth.Data;
+using Microsoft.Extensions.DependencyInjection;
+
+namespace Impostor.Server.Net.Factories
+{
+ internal class ClientFactory<TClient> : IClientFactory
+ where TClient : ClientBase
+ {
+ private int _idLast;
+ private readonly IServiceProvider _serviceProvider;
+
+ public ClientFactory(IServiceProvider serviceProvider)
+ {
+ _serviceProvider = serviceProvider;
+ }
+
+ public int NextId()
+ {
+ var clientId = Interlocked.Increment(ref _idLast);
+
+ if (clientId < 1)
+ {
+ // Super rare but reset the _idLast because of overflow.
+ _idLast = 0;
+
+ // And get a new id.
+ clientId = Interlocked.Increment(ref _idLast);
+ }
+
+ return clientId;
+ }
+
+ public async ValueTask<IClient> CreateAsync(IConnection connection, string name, int clientVersion)
+ {
+ if (clientVersion != 50516550)
+ {
+ using var packet = connection.CreateMessage(MessageType.Reliable);
+ Message01JoinGame.SerializeError(packet, false, DisconnectReason.IncorrectVersion);
+ await packet.SendAsync();
+
+ throw new ClientVersionUnsupportedException(clientVersion);
+ }
+
+ var clientId = NextId();
+ var client = ActivatorUtilities.CreateInstance<TClient>(_serviceProvider, clientId, name, connection);
+
+ await client.InitializeAsync();
+
+ return client;
+ }
+ }
+}
\ No newline at end of file
-using System.Collections.Concurrent;
-using System.Threading;
-using Hazel;
-using Impostor.Server.Exceptions;
+using System;
+using System.Collections.Concurrent;
+using System.Threading.Tasks;
+using Impostor.Server.Net.Factories;
using Microsoft.Extensions.Logging;
namespace Impostor.Server.Net.Manager
{
internal class ClientManager : IClientManager
{
- private readonly ILogger<ClientManager> _clientManager;
- private readonly GameManager _gameManager;
- private readonly ConcurrentDictionary<int, Client> _clients;
- private readonly object _idLock;
- private int _idLast;
+ private readonly ILogger<ClientManager> _logger;
+ private readonly ConcurrentDictionary<int, IClient> _clients;
+ private readonly IClientFactory _clientFactory;
- public ClientManager(ILogger<ClientManager> clientManager, GameManager gameManager)
+ public ClientManager(ILogger<ClientManager> logger, IClientFactory clientFactory)
{
- _clientManager = clientManager;
- _gameManager = gameManager;
- _clients = new ConcurrentDictionary<int, Client>();
- _idLock = new object();
- _idLast = 0;
+ _logger = logger;
+ _clientFactory = clientFactory;
+ _clients = new ConcurrentDictionary<int, IClient>();
}
- private int NextId()
+ public async ValueTask RegisterConnectionAsync(IConnection connection, string name, int clientVersion)
{
- var clientId = Interlocked.Increment(ref _idLast);
- if (clientId < 1)
+ try
{
- // Super rare but reset the _idLast because of overflow.
- _idLast = 0;
+ var client = await _clientFactory.CreateAsync(connection, name, clientVersion);
- // And get a new id.
- clientId = Interlocked.Increment(ref _idLast);
+ Register(client);
+ }
+ catch (ClientVersionUnsupportedException ex)
+ {
+ _logger.LogTrace("Closed connection because client version {Version} is not supported.", ex.Version);
}
-
- return clientId;
}
-
- public void Create(string name, Connection connection)
+
+ public void Register(IClient client)
{
- var clientId = NextId();
-
- _clientManager.LogInformation("Client connected.");
- _clients.TryAdd(clientId, new Client(this, _gameManager, clientId, name, connection));
+ _logger.LogInformation("Client connected.");
+ _clients.TryAdd(client.Id, client);
}
- public void Remove(Client client)
+ public void Remove(IClient client)
{
- _clientManager.LogInformation("Client disconnected.");
+ _logger.LogInformation("Client disconnected.");
_clients.TryRemove(client.Id, out _);
}
}
-using System.Collections.Concurrent;
+using System;
+using System.Collections.Concurrent;
using System.Collections.Generic;
using System.Linq;
using System.Net;
using Impostor.Server.Net.State;
using Impostor.Shared.Innersloth;
using Impostor.Shared.Innersloth.Data;
+using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Logging;
using Microsoft.Extensions.Options;
-using Serilog;
namespace Impostor.Server.Net.Manager
{
private readonly INodeLocator _nodeLocator;
private readonly IPEndPoint _publicIp;
private readonly ConcurrentDictionary<int, Game> _games;
+ private readonly IServiceProvider _serviceProvider;
- public GameManager(ILogger<GameManager> logger, IOptions<ServerConfig> config, INodeLocator nodeLocator)
+ public GameManager(ILogger<GameManager> logger, IOptions<ServerConfig> config, INodeLocator nodeLocator, IServiceProvider serviceProvider)
{
_logger = logger;
_nodeLocator = nodeLocator;
+ _serviceProvider = serviceProvider;
_publicIp = new IPEndPoint(IPAddress.Parse(config.Value.PublicIp), config.Value.PublicPort);
_games = new ConcurrentDictionary<int, Game>();
}
{
// TODO: Prevent duplicates when using server redirector using INodeProvider.
- var gameCode = GameCode.GenerateCode(6);
- var gameCodeStr = GameCode.IntToGameName(gameCode);
- var game = new Game(this, _nodeLocator, _publicIp, gameCode, options);
+ var gameCode = GameCode.Create();
+ var gameCodeStr = gameCode.Code;
+ var game = ActivatorUtilities.CreateInstance<Game>(_serviceProvider, _publicIp, gameCode, options);
if (_nodeLocator.Find(gameCodeStr) == null &&
_games.TryAdd(gameCode, game))
{
_nodeLocator.Save(gameCodeStr, _publicIp);
- _logger.LogDebug("Created game with code {0} ({1}).", game.CodeStr, gameCode);
+ _logger.LogDebug("Created game with code {0} ({1}).", game.Code, gameCode);
return game;
}
public void Remove(int gameCode)
{
- _logger.LogDebug("Remove game with code {0} ({1}).", GameCode.IntToGameName(gameCode), gameCode);
- _nodeLocator.Remove(GameCode.IntToGameName(gameCode));
+ _logger.LogDebug("Remove game with code {0} ({1}).", GameCodeParser.IntToGameName(gameCode), gameCode);
+ _nodeLocator.Remove(GameCodeParser.IntToGameName(gameCode));
_games.TryRemove(gameCode, out _);
}
}
+++ /dev/null
-using Hazel;
-
-namespace Impostor.Server.Net.Manager
-{
- internal interface IClientManager
- {
- void Create(string name, Connection connection);
- }
-}
\ No newline at end of file
+++ /dev/null
-using System;
-using System.Net;
-using Hazel;
-using Hazel.Udp;
-using Impostor.Server.Data;
-using Impostor.Server.Net.Manager;
-using Impostor.Server.Net.Messages;
-using Impostor.Shared.Innersloth.Data;
-using Microsoft.Extensions.Logging;
-using Microsoft.Extensions.Options;
-
-namespace Impostor.Server.Net
-{
- internal class Matchmaker
- {
- private readonly ILogger<Matchmaker> _logger;
- private readonly ServerConfig _serverConfig;
- private readonly IClientManager _clientManager;
- private readonly UdpConnectionListener _connection;
-
- public Matchmaker(
- ILogger<Matchmaker> logger,
- IOptions<ServerConfig> serverConfig,
- IClientManager clientManager)
- {
- _logger = logger;
- _serverConfig = serverConfig.Value;
- _clientManager = clientManager;
- _connection = new UdpConnectionListener(new IPEndPoint(IPAddress.Parse(_serverConfig.ListenIp), _serverConfig.ListenPort), IPMode.IPv4, s =>
- {
- _logger.LogWarning("Log from Hazel: {0}", s);
- });
-
- _connection.NewConnection += OnNewConnection;
- }
-
- public IPEndPoint EndPoint => _connection.EndPoint;
-
- private void OnNewConnection(NewConnectionEventArgs e)
- {
- try
- {
- // Handshake.
- var clientVersion = e.HandshakeData.ReadInt32();
- var clientName = e.HandshakeData.ReadString();
-
- e.HandshakeData.Recycle();
-
- if (clientVersion != 50516550)
- {
- using (var packet = MessageWriter.Get(SendOption.Reliable))
- {
- Message01JoinGame.SerializeError(packet, false, DisconnectReason.IncorrectVersion);
- e.Connection.Send(packet);
- }
- return;
- }
-
- // Create client.
- _clientManager.Create(clientName, e.Connection);
- }
- catch (Exception ex)
- {
- _logger.LogError(ex, "Error in new connection.");
- }
- }
-
- public void Start()
- {
- _connection.Start();
- }
-
- public void Stop()
- {
- _connection.Dispose();
- }
- }
-}
\ No newline at end of file
-using System.Threading;
+using System.Net;
+using System.Threading;
using System.Threading.Tasks;
using Impostor.Server.Data;
+using Impostor.Server.Net.Manager;
using Microsoft.Extensions.Hosting;
using Microsoft.Extensions.Logging;
using Microsoft.Extensions.Options;
private readonly ILogger<MatchmakerService> _logger;
private readonly ServerConfig _serverConfig;
private readonly ServerRedirectorConfig _redirectorConfig;
- private readonly Matchmaker _matchmaker;
+ private readonly IMatchmaker _matchmaker;
public MatchmakerService(
ILogger<MatchmakerService> logger,
IOptions<ServerConfig> serverConfig,
IOptions<ServerRedirectorConfig> redirectorConfig,
- Matchmaker matchmaker)
+ IMatchmaker matchmaker)
{
_logger = logger;
_serverConfig = serverConfig.Value;
_matchmaker = matchmaker;
}
- public Task StartAsync(CancellationToken cancellationToken)
+ public async Task StartAsync(CancellationToken cancellationToken)
{
- _matchmaker.Start();
+ var endpoint = new IPEndPoint(IPAddress.Parse(_serverConfig.ListenIp), _serverConfig.ListenPort);
+
+ await _matchmaker.StartAsync(endpoint);
_logger.LogInformation("Matchmaker is listening on {0}:{1}, the public server ip is {2}:{3}.",
- _matchmaker.EndPoint.Address,
- _matchmaker.EndPoint.Port,
+ endpoint.Address,
+ endpoint.Port,
_serverConfig.PublicIp,
_serverConfig.PublicPort);
? "Server redirection is enabled as master, this instance will redirect clients to other nodes."
: "Server redirection is enabled as node, this instance will accept clients.");
}
-
- return Task.CompletedTask;
}
- public Task StopAsync(CancellationToken cancellationToken)
+ public async Task StopAsync(CancellationToken cancellationToken)
{
_logger.LogWarning("Matchmaker is shutting down!");
- _matchmaker.Stop();
-
- return Task.CompletedTask;
+ await _matchmaker.StopAsync();
}
}
}
\ No newline at end of file
-using Hazel;
-using Impostor.Shared.Innersloth;
+using Impostor.Shared.Innersloth;
namespace Impostor.Server.Net.Messages
{
internal static class Message00HostGame
{
- public static void Serialize(MessageWriter writer, int gameCode)
+ public static void Serialize(IMessageWriter writer, int gameCode)
{
writer.StartMessage(MessageFlags.HostGame);
writer.Write(gameCode);
writer.EndMessage();
}
- public static GameOptionsData Deserialize(MessageReader reader)
+ public static GameOptionsData Deserialize(IMessageReader reader)
{
return GameOptionsData.Deserialize(reader.ReadBytesAndSize());
}
using System;
-using Hazel;
using Impostor.Shared.Innersloth.Data;
namespace Impostor.Server.Net.Messages
{
internal static class Message01JoinGame
{
- public static void SerializeJoin(MessageWriter writer, bool clear, int gameCode, int playerId, int hostId)
+ public static void SerializeJoin(IMessageWriter writer, bool clear, int gameCode, int playerId, int hostId)
{
if (clear)
{
- writer.Clear(SendOption.Reliable);
+ writer.Clear(MessageType.Reliable);
}
writer.StartMessage(MessageFlags.JoinGame);
writer.EndMessage();
}
- public static void SerializeError(MessageWriter writer, bool clear, DisconnectReason reason, string message = null)
+ public static void SerializeError(IMessageWriter writer, bool clear, DisconnectReason reason, string message = null)
{
if (clear)
{
- writer.Clear(SendOption.Reliable);
+ writer.Clear(MessageType.Reliable);
}
writer.StartMessage(MessageFlags.JoinGame);
writer.EndMessage();
}
- public static void Deserialize(MessageReader reader, out int gameCode, out byte unknown)
+ public static void Deserialize(IMessageReader reader, out int gameCode, out byte unknown)
{
gameCode = reader.ReadInt32();
unknown = reader.ReadByte();
-using Hazel;
-using Impostor.Shared.Innersloth.Data;
+using Impostor.Shared.Innersloth.Data;
namespace Impostor.Server.Net.Messages
{
internal static class Message04RemovePlayer
{
- public static void Serialize(MessageWriter writer, bool clear, int gameCode, int playerId, int hostId, DisconnectReason reason)
+ public static void Serialize(IMessageWriter writer, bool clear, int gameCode, int playerId, int hostId, DisconnectReason reason)
{
// Only a subset of DisconnectReason shows an unique message.
// ExitGame, Banned and Kicked.
if (clear)
{
- writer.Clear(SendOption.Reliable);
+ writer.Clear(MessageType.Reliable);
}
writer.StartMessage(MessageFlags.RemovePlayer);
writer.EndMessage();
}
- public static void Deserialize(MessageReader reader, out int playerId, out byte reason)
+ public static void Deserialize(IMessageReader reader, out int playerId, out byte reason)
{
playerId = reader.ReadPackedInt32();
reason = reader.ReadByte();
-using Hazel;
-
-namespace Impostor.Server.Net.Messages
+namespace Impostor.Server.Net.Messages
{
internal static class Message07JoinedGame
{
- public static void Serialize(MessageWriter writer, bool clear, int gameCode, int playerId, int hostId, int[] otherPlayerIds)
+ public static void Serialize(IMessageWriter writer, bool clear, int gameCode, int playerId, int hostId, int[] otherPlayerIds)
{
if (clear)
{
- writer.Clear(SendOption.Reliable);
+ writer.Clear(MessageType.Reliable);
}
writer.StartMessage(MessageFlags.JoinedGame);
-using Hazel;
-using Impostor.Shared.Innersloth.Data;
+using Impostor.Shared.Innersloth.Data;
namespace Impostor.Server.Net.Messages
{
internal static class Message10AlterGame
{
- public static void Serialize(MessageWriter writer, bool clear, int gameCode)
+ public static void Serialize(IMessageWriter writer, bool clear, int gameCode)
{
if (clear)
{
- writer.Clear(SendOption.Reliable);
+ writer.Clear(MessageType.Reliable);
}
writer.StartMessage(MessageFlags.HostGame);
writer.EndMessage();
}
- public static void Deserialize(MessageReader reader, out AlterGameTags gameTag, out bool value)
+ public static void Deserialize(IMessageReader reader, out AlterGameTags gameTag, out bool value)
{
gameTag = (AlterGameTags) reader.ReadByte();
value = reader.ReadBoolean();
-using Hazel;
-
-namespace Impostor.Server.Net.Messages
+namespace Impostor.Server.Net.Messages
{
internal static class Message11KickPlayer
{
- public static void Serialize(MessageWriter writer, bool clear, int gameCode, int playerId, bool isBan)
+ public static void Serialize(IMessageWriter writer, bool clear, int gameCode, int playerId, bool isBan)
{
if (clear)
{
- writer.Clear(SendOption.Reliable);
+ writer.Clear(MessageType.Reliable);
}
writer.StartMessage(MessageFlags.KickPlayer);
writer.EndMessage();
}
- public static void Deserialize(MessageReader reader, out int playerId, out bool isBan)
+ public static void Deserialize(IMessageReader reader, out int playerId, out bool isBan)
{
playerId = reader.ReadPackedInt32();
isBan = reader.ReadBoolean();
-using Hazel;
-
-namespace Impostor.Server.Net.Messages
+namespace Impostor.Server.Net.Messages
{
internal static class Message12WaitForHost
{
- public static void Serialize(MessageWriter writer, bool clear, int gameCode, int playerId)
+ public static void Serialize(IMessageWriter writer, bool clear, int gameCode, int playerId)
{
if (clear)
{
- writer.Clear(SendOption.Reliable);
+ writer.Clear(MessageType.Reliable);
}
writer.StartMessage(MessageFlags.WaitForHost);
using System.Net;
-using Hazel;
namespace Impostor.Server.Net.Messages
{
internal static class Message13Redirect
{
- public static void Serialize(MessageWriter writer, bool clear, IPEndPoint ipEndPoint)
+ public static void Serialize(IMessageWriter writer, bool clear, IPEndPoint ipEndPoint)
{
if (clear)
{
- writer.Clear(SendOption.Reliable);
+ writer.Clear(MessageType.Reliable);
}
writer.StartMessage(MessageFlags.Redirect);
- writer.Write(ipEndPoint.Address.GetAddressBytes());
+ writer.Write(ipEndPoint.Address);
writer.Write((ushort) ipEndPoint.Port);
writer.EndMessage();
}
using System.Collections.Generic;
-using Hazel;
using Impostor.Server.Net.State;
using Impostor.Shared.Innersloth;
{
internal static class Message16GetGameListV2
{
- public static void Deserialize(MessageReader reader, out GameOptionsData options)
+ public static void Deserialize(IMessageReader reader, out GameOptionsData options)
{
reader.ReadPackedInt32(); // Hardcoded 0.
options = GameOptionsData.Deserialize(reader.ReadBytesAndSize());
}
- public static void Serialize(MessageWriter writer, int skeldGameCount, int miraHqGameCount, int polusGameCount, IEnumerable<Game> games)
+ public static void Serialize(IMessageWriter writer, int skeldGameCount, int miraHqGameCount, int polusGameCount, IEnumerable<Game> games)
{
writer.StartMessage(MessageFlags.GetGameListV2);
foreach (var game in games)
{
writer.StartMessage(0);
- writer.Write(game.PublicIp.Address.GetAddressBytes());
+ writer.Write(game.PublicIp.Address);
writer.Write((ushort) game.PublicIp.Port);
writer.Write(game.Code);
writer.Write(game.Host.Client.Name);
+++ /dev/null
-using System.Collections.Generic;
-using Hazel;
-using Impostor.Server.Net.Manager;
-using Microsoft.Extensions.Logging;
-
-namespace Impostor.Server.Net.Redirector
-{
- internal class ClientManagerRedirector : IClientManager
- {
- private readonly ILogger<ClientManagerRedirector> _logger;
- private readonly INodeProvider _nodeProvider;
- private readonly INodeLocator _nodeLocator;
- private readonly HashSet<ClientRedirector> _clients;
-
- public ClientManagerRedirector(ILogger<ClientManagerRedirector> logger, INodeProvider nodeProvider, INodeLocator nodeLocator)
- {
- _logger = logger;
- _nodeProvider = nodeProvider;
- _nodeLocator = nodeLocator;
- _clients = new HashSet<ClientRedirector>();
- }
-
- public void Create(string name, Connection connection)
- {
- _logger.LogInformation("Client connected.");
- _clients.Add(new ClientRedirector(name, connection, this, _nodeProvider, _nodeLocator));
- }
-
- public void Remove(ClientRedirector client)
- {
- _logger.LogInformation("Client disconnected.");
- _clients.Remove(client);
- }
- }
-}
\ No newline at end of file
-using System;
-using Hazel;
+using System.Threading.Tasks;
using Impostor.Server.Data;
+using Impostor.Server.Net.Manager;
using Impostor.Server.Net.Messages;
using Impostor.Shared.Innersloth;
using Impostor.Shared.Innersloth.Data;
namespace Impostor.Server.Net.Redirector
{
- internal class ClientRedirector
+ internal class ClientRedirector : ClientBase
{
private static readonly ILogger Logger = Log.ForContext<ClientRedirector>();
- private readonly string _name;
- private readonly Connection _connection;
- private readonly ClientManagerRedirector _clientManager;
+ private readonly IClientManager _clientManager;
private readonly INodeProvider _nodeProvider;
private readonly INodeLocator _nodeLocator;
- public ClientRedirector(string name, Connection connection, ClientManagerRedirector clientManager, INodeProvider nodeProvider, INodeLocator nodeLocator)
+ public ClientRedirector(int id, string name, IConnection connection, IClientManager clientManager, INodeProvider nodeProvider, INodeLocator nodeLocator)
+ : base(id, name, connection)
{
- _name = name;
- _connection = connection;
- _connection.DataReceived += OnDataReceived;
- _connection.Disconnected += OnDisconnected;
_clientManager = clientManager;
_nodeProvider = nodeProvider;
_nodeLocator = nodeLocator;
}
- private void OnDataReceived(DataReceivedEventArgs e)
+ protected override async ValueTask OnMessageReceived(IMessage message)
{
- try
- {
- while (true)
- {
- if (e.Message.Position >= e.Message.Length)
- {
- break;
- }
-
- OnMessageReceived(e.Message.ReadMessage());
- }
- }
- catch (Exception ex)
- {
- Logger.Error(ex, "Exception caught in client data handler.");
- }
- }
-
- private void OnMessageReceived(MessageReader message)
- {
- var flag = message.Tag;
+ var reader = message.CreateReader();
+ var flag = reader.Tag;
Logger.Verbose("Server got {0}.", flag);
{
case MessageFlags.HostGame:
{
- using (var packet = MessageWriter.Get(SendOption.Reliable))
- {
- Message13Redirect.Serialize(packet, false, _nodeProvider.Get());
- _connection.Send(packet);
- }
+ using var packet = Connection.CreateMessage(MessageType.Reliable);
+ Message13Redirect.Serialize(packet, false, _nodeProvider.Get());
+ await packet.SendAsync();
break;
}
case MessageFlags.JoinGame:
{
- Message01JoinGame.Deserialize(message,
+ Message01JoinGame.Deserialize(reader,
out var gameCode,
out var unknown);
- using (var packet = MessageWriter.Get(SendOption.Reliable))
+ using (var packet = Connection.CreateMessage(MessageType.Reliable))
{
- var endpoint = _nodeLocator.Find(GameCode.IntToGameName(gameCode));
+ var endpoint = _nodeLocator.Find(GameCodeParser.IntToGameName(gameCode));
if (endpoint == null)
{
Message01JoinGame.SerializeError(packet, false, DisconnectReason.GameMissing);
{
Message13Redirect.Serialize(packet, false, endpoint);
}
-
- _connection.Send(packet);
+
+ await packet.SendAsync();
}
break;
}
case MessageFlags.GetGameListV2:
{
// TODO: Implement.
- using (var packet = MessageWriter.Get(SendOption.Reliable))
+ using (var packet = Connection.CreateMessage(MessageType.Reliable))
{
Message01JoinGame.SerializeError(packet, false, DisconnectReason.Custom, DisconnectMessages.NotImplemented);
- _connection.Send(packet);
+ await packet.SendAsync();
}
break;
}
}
}
- private void OnDisconnected(object sender, DisconnectedEventArgs e)
+ protected override ValueTask OnDisconnected()
{
_clientManager.Remove(this);
+ return default;
}
}
}
\ No newline at end of file
-using Hazel;
+using System.Threading.Tasks;
using Impostor.Server.Net.Messages;
using Impostor.Shared.Innersloth;
using Impostor.Shared.Innersloth.Data;
/// All options given.
/// At this moment, the client can only specify the map, impostor count and chat language.
/// </param>
- public void OnRequestGameList(GameOptionsData options)
+ public async ValueTask OnRequestGameList(GameOptionsData options)
{
- using (var message = MessageWriter.Get(SendOption.Reliable))
+ using (var message = Client.Connection.CreateMessage(MessageType.Reliable))
{
var games = _gameManager.FindListings((MapFlags) options.MapId, options.NumImpostors, options.Keywords);
var polusGameCount = _gameManager.GetGameCount(MapFlags.Polus);
Message16GetGameListV2.Serialize(message, skeldGameCount, miraHqGameCount, polusGameCount, games);
-
- Client.Send(message);
+
+ await message.SendAsync();
}
}
}
-using Hazel;
+using System.Threading.Tasks;
using Impostor.Server.Net.Manager;
using Impostor.Server.Net.Messages;
using Impostor.Shared.Innersloth.Data;
namespace Impostor.Server.Net.State
{
- internal partial class ClientPlayer
+ internal partial class ClientPlayer : IClientPlayer
{
private readonly GameManager _gameManager;
public Game Game { get; set; }
public LimboStates Limbo { get; set; }
- public void SendDisconnectReason(DisconnectReason reason, string message = null)
+ public async ValueTask SendDisconnectReason(DisconnectReason reason, string message = null)
{
- using (var packet = MessageWriter.Get(SendOption.Reliable))
+ using (var packet = Client.Connection.CreateMessage(MessageType.Reliable))
{
Message01JoinGame.SerializeError(packet, false, reason, message);
- Client.Connection.Send(packet);
+ await packet.SendAsync();
}
}
+
+ IClient IClientPlayer.Client => Client;
+
+ IGame IClientPlayer.Game => Game;
}
}
\ No newline at end of file
using System;
-using System.Linq;
-using Hazel;
+using System.Threading.Tasks;
using Impostor.Server.Data;
-using Impostor.Server.Exceptions;
using Impostor.Shared.Innersloth.Data;
namespace Impostor.Server.Net.State
{
internal partial class Game
{
- public void HandleStartGame(MessageReader message)
+ public async ValueTask HandleStartGame(IMessageReader message)
{
GameState = GameStates.Started;
-
- using (var packet = MessageWriter.Get(SendOption.Reliable))
- {
- packet.CopyFrom(message);
- SendToAllExcept(packet, null);
- }
+
+ using var packet = CreateMessage(MessageType.Reliable);
+ message.CopyTo(packet);
+ await packet.SendToAllAsync(LimboStates.NotLimbo);
}
- public void HandleJoinGame(ClientPlayer sender)
+ public async ValueTask HandleJoinGame(ClientPlayer sender)
{
// Check if the IP of the player is banned.
if (_bannedIps.Contains(sender.Client.Connection.EndPoint.Address))
{
- sender.SendDisconnectReason(DisconnectReason.Banned);
+ await sender.SendDisconnectReason(DisconnectReason.Banned);
return;
}
// - The game is full.
if (sender.Game != this && _players.Count >= Options.MaxPlayers)
{
- sender.SendDisconnectReason(DisconnectReason.GameFull);
+ await sender.SendDisconnectReason(DisconnectReason.GameFull);
return;
}
// Check current player state.
if (sender.Limbo == LimboStates.NotLimbo)
{
- sender.SendDisconnectReason(DisconnectReason.Custom, "Invalid limbo state while joining.");
+ await sender.SendDisconnectReason(DisconnectReason.Custom, "Invalid limbo state while joining.");
return;
}
switch (GameState)
{
case GameStates.NotStarted:
- HandleJoinGameNew(sender);
+ await HandleJoinGameNew(sender);
break;
case GameStates.Ended:
- HandleJoinGameNext(sender);
+ await HandleJoinGameNext(sender);
break;
case GameStates.Started:
- sender.SendDisconnectReason(DisconnectReason.GameStarted);
+ await sender.SendDisconnectReason(DisconnectReason.GameStarted);
return;
case GameStates.Destroyed:
- sender.SendDisconnectReason(DisconnectReason.Custom, DisconnectMessages.Destroyed);
+ await sender.SendDisconnectReason(DisconnectReason.Custom, DisconnectMessages.Destroyed);
return;
default:
throw new ArgumentOutOfRangeException();
}
}
- public void HandleEndGame(MessageReader message)
+ public async ValueTask HandleEndGame(IMessageReader message)
{
GameState = GameStates.Ended;
// Broadcast end of the game.
- using (var packet = MessageWriter.Get(SendOption.Reliable))
+ using (var packet = CreateMessage(MessageType.Reliable))
{
- packet.CopyFrom(message);
- SendToAllExcept(packet, null);
+ message.CopyTo(packet);
+ await packet.SendToAllAsync(LimboStates.NotLimbo);
}
// Put all players in the correct limbo state.
}
}
- public void HandleAlterGame(MessageReader message, ClientPlayer sender, bool isPublic)
+ public async ValueTask HandleAlterGame(IMessageReader message, ClientPlayer sender, bool isPublic)
{
IsPublic = isPublic;
-
- using (var packet = MessageWriter.Get(SendOption.Reliable))
- {
- packet.CopyFrom(message);
- SendToAllExcept(packet, sender.Client.Id);
- }
+
+ using var packet = CreateMessage(MessageType.Reliable);
+ message.CopyTo(packet);
+ await packet.SendToAllExceptAsync(LimboStates.NotLimbo, sender.Client.Id);
}
- public void HandleRemovePlayer(int playerId, DisconnectReason reason)
+ public async ValueTask HandleRemovePlayer(int playerId, DisconnectReason reason)
{
- PlayerRemove(playerId, out _);
+ await PlayerRemove(playerId);
// It's possible that the last player was removed, so check if the game is still around.
if (GameState == GameStates.Destroyed)
{
return;
}
-
- using (var packet = MessageWriter.Get(SendOption.Reliable))
- {
- WriteRemovePlayerMessage(packet, false, playerId, reason);
- SendToAllExcept(packet, playerId);
- }
+
+ using var packet = CreateMessage(MessageType.Reliable);
+ WriteRemovePlayerMessage(packet, false, playerId, reason);
+ await packet.SendToAllExceptAsync(LimboStates.NotLimbo, playerId);
}
- public void HandleKickPlayer(int playerId, bool isBan)
+ public async ValueTask HandleKickPlayer(int playerId, bool isBan)
{
- Logger.Information("{0} - Player {1} has left.", CodeStr, playerId);
+ Logger.Information("{0} - Player {1} has left.", Code, playerId);
- using (var message = MessageWriter.Get(SendOption.Reliable))
+ using (var message = CreateMessage(MessageType.Reliable))
{
// Send message to everyone that this player was kicked.
WriteKickPlayerMessage(message, false, playerId, isBan);
- SendToAllExcept(message, null);
-
- if (PlayerRemove(playerId, out var player) && isBan)
- {
- _bannedIps.Add(player.Client.Connection.EndPoint.Address);
- }
+ await message.SendToAllAsync(LimboStates.NotLimbo);
+
+ await PlayerRemove(playerId, isBan);
// Rmeove the player from everyone's game.
WriteRemovePlayerMessage(message, true, playerId, isBan
? DisconnectReason.Banned
: DisconnectReason.Kicked);
- SendToAllExcept(message, player?.Client.Id);
+ await message.SendToAllExceptAsync(LimboStates.NotLimbo, playerId);
}
}
- private void HandleJoinGameNew(ClientPlayer sender)
+ private async ValueTask HandleJoinGameNew(ClientPlayer sender)
{
- Logger.Information("{0} - Player {1} ({2}) is joining.", CodeStr, sender.Client.Name, sender.Client.Id);
+ Logger.Information("{0} - Player {1} ({2}) is joining.", Code, sender.Client.Name, sender.Client.Id);
// Add player to the game.
if (sender.Game == null)
PlayerAdd(sender);
}
- using (var message = MessageWriter.Get(SendOption.Reliable))
+ using (var message = CreateMessage(MessageType.Reliable))
{
WriteJoinedGameMessage(message, false, sender);
WriteAlterGameMessage(message, false);
sender.Limbo = LimboStates.NotLimbo;
- sender.Client.Send(message);
+ await message.SendToAsync(sender.Client);
- BroadcastJoinMessage(message, true, sender);
+ await BroadcastJoinMessage(message, true, sender);
}
}
- private void HandleJoinGameNext(ClientPlayer sender)
+ private async ValueTask HandleJoinGameNext(ClientPlayer sender)
{
- Logger.Information("{0} - Player {1} ({2}) is rejoining.", CodeStr, sender.Client.Name, sender.Client.Id);
+ Logger.Information("{0} - Player {1} ({2}) is rejoining.", Code, sender.Client.Name, sender.Client.Id);
// Add player to the game.
if (sender.Game == null)
GameState = GameStates.NotStarted;
// Spawn the host.
- HandleJoinGameNew(sender);
+ await HandleJoinGameNew(sender);
// Pull players out of limbo.
- CheckLimboPlayers();
+ await CheckLimboPlayers();
return;
}
sender.Limbo = LimboStates.WaitingForHost;
- using (var packet = MessageWriter.Get(SendOption.Reliable))
- {
- WriteWaitForHostMessage(packet, false, sender);
- sender.Client.Send(packet);
+ using var packet = CreateMessage(MessageType.Reliable);
+
+ WriteWaitForHostMessage(packet, false, sender);
+ await packet.SendToAsync(sender.Client);
- BroadcastJoinMessage(packet, true, sender);
- }
+ await BroadcastJoinMessage(packet, true, sender);
}
}
}
\ No newline at end of file
using System.Linq;
-using Hazel;
using Impostor.Server.Net.Messages;
using Impostor.Shared.Innersloth.Data;
{
internal partial class Game
{
- private void WriteRemovePlayerMessage(MessageWriter message, bool clear, int playerId, DisconnectReason reason)
+ private void WriteRemovePlayerMessage(IMessageWriter message, bool clear, int playerId, DisconnectReason reason)
{
Message04RemovePlayer.Serialize(message, clear, Code, playerId, HostId, reason);
}
- private void WriteJoinedGameMessage(MessageWriter message, bool clear, ClientPlayer player)
+ private void WriteJoinedGameMessage(IMessageWriter message, bool clear, ClientPlayer player)
{
var playerIds = _players
.Where(x => x.Value != player)
Message07JoinedGame.Serialize(message, clear, Code, player.Client.Id, HostId, playerIds);
}
- private void WriteAlterGameMessage(MessageWriter message, bool clear)
+ private void WriteAlterGameMessage(IMessageWriter message, bool clear)
{
Message10AlterGame.Serialize(message, clear, Code);
}
- private void WriteKickPlayerMessage(MessageWriter message, bool clear, int playerId, bool isBan)
+ private void WriteKickPlayerMessage(IMessageWriter message, bool clear, int playerId, bool isBan)
{
Message11KickPlayer.Serialize(message, clear, Code, playerId, isBan);
}
- private void WriteWaitForHostMessage(MessageWriter message, bool clear, ClientPlayer player)
+ private void WriteWaitForHostMessage(IMessageWriter message, bool clear, ClientPlayer player)
{
Message12WaitForHost.Serialize(message, clear, Code, player.Client.Id);
}
-using System.Linq;
-using Hazel;
+using System.Collections.Generic;
+using System.Linq;
+using System.Threading.Tasks;
using Impostor.Server.Exceptions;
using Impostor.Shared.Innersloth.Data;
}
}
- private bool PlayerRemove(int playerId, out ClientPlayer player)
+ private async ValueTask<bool> PlayerRemove(int playerId, bool isBan = false)
{
- if (!_players.TryRemove(playerId, out player))
+ if (!_players.TryRemove(playerId, out var player))
{
return false;
}
player.Limbo = LimboStates.PreSpawn;
player.Game = null;
- Logger.Information("{0} - Player {1} ({2}) has left.", CodeStr, player.Client.Name, playerId);
+ Logger.Information("{0} - Player {1} ({2}) has left.", Code, player.Client.Name, playerId);
// Game is empty, remove it.
if (_players.Count == 0)
// Host migration.
if (HostId == playerId)
{
- MigrateHost();
+ await MigrateHost();
+ }
+
+ if (isBan)
+ {
+ _bannedIps.Add(player.Client.Connection.EndPoint.Address);
}
return true;
}
- private void MigrateHost()
+ private async ValueTask MigrateHost()
{
// Pick the first player as new host.
var host = _players.First().Value;
HostId = host.Client.Id;
- Logger.Information("{0} - Assigned {1} ({2}) as new host.", CodeStr, host.Client.Name, host.Client.Id);
+ Logger.Information("{0} - Assigned {1} ({2}) as new host.", Code, host.Client.Name, host.Client.Id);
// Check our current game state.
if (GameState == GameStates.Ended && host.Limbo == LimboStates.WaitingForHost)
GameState = GameStates.NotStarted;
// Spawn the host.
- HandleJoinGameNew(host);
+ await HandleJoinGameNew(host);
// Pull players out of limbo.
- CheckLimboPlayers();
+ await CheckLimboPlayers();
}
}
- private void CheckLimboPlayers()
+ private async ValueTask CheckLimboPlayers()
{
- using (var message = MessageWriter.Get(SendOption.Reliable))
+ using var message = CreateMessage(MessageType.Reliable);
+
+ foreach (var (_, player) in _players.Where(x => x.Value.Limbo == LimboStates.WaitingForHost))
{
- foreach (var (_, player) in _players.Where(x => x.Value.Limbo == LimboStates.WaitingForHost))
- {
- WriteJoinedGameMessage(message, true, player);
- WriteAlterGameMessage(message, false);
+ WriteJoinedGameMessage(message, true, player);
+ WriteAlterGameMessage(message, false);
- player.Limbo = LimboStates.NotLimbo;
- player.Client.Send(message);
- }
+ player.Limbo = LimboStates.NotLimbo;
+ await message.SendToAsync(player.Client);
}
}
}
using System.Collections.Generic;
using System.Linq;
using System.Net;
-using Hazel;
+using System.Threading.Tasks;
using Impostor.Server.Net.Manager;
using Impostor.Server.Net.Messages;
using Impostor.Server.Net.Redirector;
namespace Impostor.Server.Net.State
{
- internal partial class Game
+ internal partial class Game : IGame
{
private static readonly ILogger Logger = Log.ForContext<Game>();
private readonly GameManager _gameManager;
private readonly INodeLocator _nodeLocator;
+ private readonly IMatchmaker matchmaker;
private readonly ConcurrentDictionary<int, ClientPlayer> _players;
private readonly HashSet<IPAddress> _bannedIps;
- public Game(GameManager gameManager, INodeLocator nodeLocator, IPEndPoint publicIp, int code, GameOptionsData options)
+ public Game(
+ GameManager gameManager,
+ INodeLocator nodeLocator,
+ IPEndPoint publicIp,
+ GameCode code,
+ GameOptionsData options,
+ IMatchmaker matchmaker)
{
_gameManager = gameManager;
_nodeLocator = nodeLocator;
PublicIp = publicIp;
Code = code;
- CodeStr = GameCode.IntToGameName(code);
HostId = -1;
GameState = GameStates.NotStarted;
Options = options;
+ this.matchmaker = matchmaker;
+ Items = new ConcurrentDictionary<object, object>();
}
public IPEndPoint PublicIp { get; }
- public int Code { get; }
- public string CodeStr { get; }
+ public GameCode Code { get; }
public bool IsPublic { get; private set; }
public int HostId { get; private set; }
public GameStates GameState { get; private set; }
public GameOptionsData Options { get; }
+ public IDictionary<object, object> Items { get; }
public int PlayerCount => _players.Count;
- public ClientPlayer Host => _players[HostId];
- /// <summary>
- /// Send a message to all players except one.
- /// </summary>
- /// <param name="message">The message to send.</param>
- /// <param name="senderId">
- /// The player to exclude from sending the message.
- /// Set to null to send a message to everyone.
- /// </param>
- public void SendToAllExcept(MessageWriter message, int? senderId)
+ public IClientPlayer Host => _players[HostId];
+
+ private ValueTask BroadcastJoinMessage(IGameMessageWriter message, bool clear, ClientPlayer player)
{
- foreach (var (_, player) in _players.Where(x =>
- x.Value.Limbo == LimboStates.NotLimbo &&
- x.Value.Client.Id != senderId))
- {
- if (player.Client.Connection.State != ConnectionState.Connected)
- {
- Logger.Warning("[{0}] Tried to send data to a disconnected player ({1}).", senderId, player.Client.Id);
- continue;
- }
-
- player.Client.Send(message);
- }
+ Message01JoinGame.SerializeJoin(message, clear, Code, player.Client.Id, HostId);
+
+ return message.SendToAllExceptAsync(LimboStates.NotLimbo, player.Client.Id);
}
- /// <summary>
- /// Send a message to a specific player.
- /// </summary>
- /// <param name="message">The message to send.</param>
- /// <param name="playerId"></param>
- public void SendTo(MessageWriter message, int playerId)
+ public IEnumerable<IClientPlayer> Players => _players.Select(p => p.Value);
+
+ public IGameMessageWriter CreateMessage(MessageType type)
{
- if (_players.TryGetValue(playerId, out var player))
- {
- if (player.Client.Connection.State != ConnectionState.Connected)
- {
- Logger.Warning("[{0}] Sending data to {1} failed, player is not connected.", CodeStr, player.Client.Id);
- return;
- }
-
- player.Client.Send(message);
- }
- else
- {
- Logger.Warning("[{0}] Sending data to {1} failed, player does not exist.", CodeStr, playerId);
- }
+ return matchmaker.CreateGameMessageWriter(this, type);
}
-
- private void BroadcastJoinMessage(MessageWriter message, bool clear, ClientPlayer player)
+
+ public bool TryGetPlayer(int id, out IClientPlayer player)
{
- Message01JoinGame.SerializeJoin(message, clear, Code, player.Client.Id, HostId);
-
- SendToAllExcept(message, player.Client.Id);
+ if (_players.TryGetValue(id, out var result))
+ {
+ player = result;
+ return true;
+ }
+
+ player = default;
+ return false;
}
}
}
\ No newline at end of file
using System;
-using System.IO;
using Impostor.Server.Data;
+using Impostor.Server.Hazel;
using Impostor.Server.Net;
+using Impostor.Server.Net.Factories;
using Impostor.Server.Net.Manager;
using Impostor.Server.Net.Redirector;
using Microsoft.Extensions.Configuration;
services.AddSingleton<INodeLocator, NodeLocatorNoOp>();
}
+ services.AddSingleton<IClientManager, ClientManager>();
+
if (redirector.Enabled && redirector.Master)
{
- services.AddSingleton<IClientManager, ClientManagerRedirector>();
+ services.AddSingleton<IClientFactory, ClientFactory<ClientRedirector>>();
// For a master server, we don't need a GameManager.
}
else
{
- services.AddSingleton<IClientManager, ClientManager>();
+ services.AddSingleton<IClientFactory, ClientFactory<Client>>();
services.AddSingleton<GameManager>();
}
-
- services.AddSingleton<Matchmaker>();
+
+ services.UseHazelMatchmaking();
services.AddHostedService<MatchmakerService>();
})
.UseConsoleLifetime()
+++ /dev/null
-namespace Impostor.Shared.Innersloth.Data
-{
- public enum LimboStates
- {
- PreSpawn,
- NotLimbo,
- WaitingForHost,
- }
-}
\ No newline at end of file
+++ /dev/null
-using System;
-using System.Buffers.Binary;
-using System.Linq;
-using System.Security.Cryptography;
-using System.Text;
-
-namespace Impostor.Shared.Innersloth
-{
- public static class GameCode
- {
- private const string V2 = "QWXRTYLPESDFGHUJKZOCVBINMA";
- private static readonly int[] V2Map = {
- 25,
- 21,
- 19,
- 10,
- 8,
- 11,
- 12,
- 13,
- 22,
- 15,
- 16,
- 6,
- 24,
- 23,
- 18,
- 7,
- 0,
- 3,
- 9,
- 4,
- 14,
- 20,
- 1,
- 2,
- 5,
- 17
- };
- private static readonly RNGCryptoServiceProvider Random = new RNGCryptoServiceProvider();
-
- public static string IntToGameName(int input)
- {
- // V2.
- if (input < -1)
- {
- return IntToGameNameV2(input);
- }
-
- // V1.
- Span<byte> code = stackalloc byte[4];
- BinaryPrimitives.WriteInt32LittleEndian(code, input);
-#if NET452
- return Encoding.UTF8.GetString(code.Slice(0, 4).ToArray());
-#else
- return Encoding.UTF8.GetString(code.Slice(0, 4));
-#endif
- }
-
- private static string IntToGameNameV2(int input)
- {
- var a = input & 0x3FF;
- var b = (input >> 10) & 0xFFFFF;
-
- return new string(new []
- {
- V2[a % 26],
- V2[a / 26],
- V2[b % 26],
- V2[b / 26 % 26],
- V2[b / (26 * 26) % 26],
- V2[b / (26 * 26 * 26) % 26]
- });
- }
-
- public static int GameNameToInt(string code)
- {
- var upper = code.ToUpperInvariant();
- if (upper.Any(x => !char.IsLetter(x)))
- {
- return -1;
- }
-
- var len = code.Length;
- if (len == 6)
- {
- return GameNameToIntV2(upper);
- }
-
- if (len == 4)
- {
- return code[0] | ((code[1] | ((code[2] | (code[3] << 8)) << 8)) << 8);
- }
-
- return -1;
- }
-
- private static int GameNameToIntV2(string code)
- {
- var a = V2Map[code[0] - 65];
- var b = V2Map[code[1] - 65];
- var c = V2Map[code[2] - 65];
- var d = V2Map[code[3] - 65];
- var e = V2Map[code[4] - 65];
- var f = V2Map[code[5] - 65];
-
- var one = (a + 26 * b) & 0x3FF;
- var two = (c + 26 * (d + 26 * (e + 26 * f)));
-
- return (int) (one | ((two << 10) & 0x3FFFFC00) | 0x80000000);
- }
-
- public static int GenerateCode(int len)
- {
- if (len != 4 && len != 6)
- {
- throw new ArgumentException("should be 4 or 6", nameof(len));
- }
-
- // Generate random bytes.
-#if NET452
- var data = new byte[len];
-#else
- Span<byte> data = stackalloc byte[len];
-#endif
- Random.GetBytes(data);
-
- // Convert to their char representation.
- Span<char> dataChar = stackalloc char[len];
- for (var i = 0; i < len; i++)
- {
- dataChar[i] = V2[V2Map[data[i] % 26]];
- }
-
-#if NET452
- return GameNameToInt(new string(dataChar.ToArray()));
-#else
- return GameNameToInt(new string(dataChar));
-#endif
- }
- }
-}
\ No newline at end of file
--- /dev/null
+using System;
+using System.Buffers.Binary;
+using System.Linq;
+using System.Security.Cryptography;
+using System.Text;
+
+namespace Impostor.Shared.Innersloth
+{
+ public static class GameCodeParser
+ {
+ private const string V2 = "QWXRTYLPESDFGHUJKZOCVBINMA";
+ private static readonly int[] V2Map = {
+ 25,
+ 21,
+ 19,
+ 10,
+ 8,
+ 11,
+ 12,
+ 13,
+ 22,
+ 15,
+ 16,
+ 6,
+ 24,
+ 23,
+ 18,
+ 7,
+ 0,
+ 3,
+ 9,
+ 4,
+ 14,
+ 20,
+ 1,
+ 2,
+ 5,
+ 17
+ };
+ private static readonly RNGCryptoServiceProvider Random = new RNGCryptoServiceProvider();
+
+ public static string IntToGameName(int input)
+ {
+ // V2.
+ if (input < -1)
+ {
+ return IntToGameNameV2(input);
+ }
+
+ // V1.
+ Span<byte> code = stackalloc byte[4];
+ BinaryPrimitives.WriteInt32LittleEndian(code, input);
+#if NET452
+ return Encoding.UTF8.GetString(code.Slice(0, 4).ToArray());
+#else
+ return Encoding.UTF8.GetString(code.Slice(0, 4));
+#endif
+ }
+
+ private static string IntToGameNameV2(int input)
+ {
+ var a = input & 0x3FF;
+ var b = (input >> 10) & 0xFFFFF;
+
+ return new string(new []
+ {
+ V2[a % 26],
+ V2[a / 26],
+ V2[b % 26],
+ V2[b / 26 % 26],
+ V2[b / (26 * 26) % 26],
+ V2[b / (26 * 26 * 26) % 26]
+ });
+ }
+
+ public static int GameNameToInt(string code)
+ {
+ var upper = code.ToUpperInvariant();
+ if (upper.Any(x => !char.IsLetter(x)))
+ {
+ return -1;
+ }
+
+ var len = code.Length;
+ if (len == 6)
+ {
+ return GameNameToIntV2(upper);
+ }
+
+ if (len == 4)
+ {
+ return code[0] | ((code[1] | ((code[2] | (code[3] << 8)) << 8)) << 8);
+ }
+
+ return -1;
+ }
+
+ private static int GameNameToIntV2(string code)
+ {
+ var a = V2Map[code[0] - 65];
+ var b = V2Map[code[1] - 65];
+ var c = V2Map[code[2] - 65];
+ var d = V2Map[code[3] - 65];
+ var e = V2Map[code[4] - 65];
+ var f = V2Map[code[5] - 65];
+
+ var one = (a + 26 * b) & 0x3FF;
+ var two = (c + 26 * (d + 26 * (e + 26 * f)));
+
+ return (int) (one | ((two << 10) & 0x3FFFFC00) | 0x80000000);
+ }
+
+ public static int GenerateCode(int len)
+ {
+ if (len != 4 && len != 6)
+ {
+ throw new ArgumentException("should be 4 or 6", nameof(len));
+ }
+
+ // Generate random bytes.
+#if NET452
+ var data = new byte[len];
+#else
+ Span<byte> data = stackalloc byte[len];
+#endif
+ Random.GetBytes(data);
+
+ // Convert to their char representation.
+ Span<char> dataChar = stackalloc char[len];
+ for (var i = 0; i < len; i++)
+ {
+ dataChar[i] = V2[V2Map[data[i] % 26]];
+ }
+
+#if NET452
+ return GameNameToInt(new string(dataChar.ToArray()));
+#else
+ return GameNameToInt(new string(dataChar));
+#endif
+ }
+ }
+}
\ No newline at end of file
}
}
- public static GameOptionsData Deserialize(byte[] bytes)
+ public static GameOptionsData Deserialize(ReadOnlyMemory<byte> bytes)
{
- using (var stream = new MemoryStream(bytes))
+ // TODO: Remove memory allocation.
+
+ using (var stream = new MemoryStream(bytes.ToArray()))
using (var reader = new BinaryReader(stream))
{
var result = new GameOptionsData();
const string code = "ABCD";
const int codeInt = 0x44434241;
- Assert.Equal(code, GameCode.IntToGameName(codeInt));
- Assert.Equal(codeInt, GameCode.GameNameToInt(code));
+ Assert.Equal(code, GameCodeParser.IntToGameName(codeInt));
+ Assert.Equal(codeInt, GameCodeParser.GameNameToInt(code));
}
[Fact]
const string code = "ABCDEF";
const int codeInt = -1943683525;
- Assert.Equal(code, GameCode.IntToGameName(codeInt));
- Assert.Equal(codeInt, GameCode.GameNameToInt(code));
+ Assert.Equal(code, GameCodeParser.IntToGameName(codeInt));
+ Assert.Equal(codeInt, GameCodeParser.GameNameToInt(code));
}
}
}
\ No newline at end of file
EndProject
Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Impostor.Client", "Impostor.Client\Impostor.Client.csproj", "{804CF172-0C87-4423-9688-BD97D549891E}"
EndProject
+Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Impostor.Server.Api", "Impostor.Server.Api\Impostor.Server.Api.csproj", "{E096A7D7-D693-4A13-A526-38CC574D84F8}"
+EndProject
+Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Imposter.Reactive", "Imposter.Reactive\Imposter.Reactive.csproj", "{7D0541DF-5BD8-4175-BAF1-14278B7894F1}"
+EndProject
+Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Impostor.Server.Hazel", "Impostor.Server.Hazel\Impostor.Server.Hazel.csproj", "{C9E8E1E2-BFE3-41AD-8CB9-82F0F52E2306}"
+EndProject
Global
GlobalSection(SolutionConfigurationPlatforms) = preSolution
Debug|Any CPU = Debug|Any CPU
{804CF172-0C87-4423-9688-BD97D549891E}.Release|Any CPU.Build.0 = Release|Any CPU
{804CF172-0C87-4423-9688-BD97D549891E}.Release|x86.ActiveCfg = Release|Any CPU
{804CF172-0C87-4423-9688-BD97D549891E}.Release|x86.Build.0 = Release|Any CPU
+ {E096A7D7-D693-4A13-A526-38CC574D84F8}.Debug|Any CPU.ActiveCfg = Debug|Any CPU
+ {E096A7D7-D693-4A13-A526-38CC574D84F8}.Debug|Any CPU.Build.0 = Debug|Any CPU
+ {E096A7D7-D693-4A13-A526-38CC574D84F8}.Debug|x86.ActiveCfg = Debug|Any CPU
+ {E096A7D7-D693-4A13-A526-38CC574D84F8}.Debug|x86.Build.0 = Debug|Any CPU
+ {E096A7D7-D693-4A13-A526-38CC574D84F8}.Release|Any CPU.ActiveCfg = Release|Any CPU
+ {E096A7D7-D693-4A13-A526-38CC574D84F8}.Release|Any CPU.Build.0 = Release|Any CPU
+ {E096A7D7-D693-4A13-A526-38CC574D84F8}.Release|x86.ActiveCfg = Release|Any CPU
+ {E096A7D7-D693-4A13-A526-38CC574D84F8}.Release|x86.Build.0 = Release|Any CPU
+ {7D0541DF-5BD8-4175-BAF1-14278B7894F1}.Debug|Any CPU.ActiveCfg = Debug|Any CPU
+ {7D0541DF-5BD8-4175-BAF1-14278B7894F1}.Debug|Any CPU.Build.0 = Debug|Any CPU
+ {7D0541DF-5BD8-4175-BAF1-14278B7894F1}.Debug|x86.ActiveCfg = Debug|Any CPU
+ {7D0541DF-5BD8-4175-BAF1-14278B7894F1}.Debug|x86.Build.0 = Debug|Any CPU
+ {7D0541DF-5BD8-4175-BAF1-14278B7894F1}.Release|Any CPU.ActiveCfg = Release|Any CPU
+ {7D0541DF-5BD8-4175-BAF1-14278B7894F1}.Release|Any CPU.Build.0 = Release|Any CPU
+ {7D0541DF-5BD8-4175-BAF1-14278B7894F1}.Release|x86.ActiveCfg = Release|Any CPU
+ {7D0541DF-5BD8-4175-BAF1-14278B7894F1}.Release|x86.Build.0 = Release|Any CPU
+ {C9E8E1E2-BFE3-41AD-8CB9-82F0F52E2306}.Debug|Any CPU.ActiveCfg = Debug|Any CPU
+ {C9E8E1E2-BFE3-41AD-8CB9-82F0F52E2306}.Debug|Any CPU.Build.0 = Debug|Any CPU
+ {C9E8E1E2-BFE3-41AD-8CB9-82F0F52E2306}.Debug|x86.ActiveCfg = Debug|Any CPU
+ {C9E8E1E2-BFE3-41AD-8CB9-82F0F52E2306}.Debug|x86.Build.0 = Debug|Any CPU
+ {C9E8E1E2-BFE3-41AD-8CB9-82F0F52E2306}.Release|Any CPU.ActiveCfg = Release|Any CPU
+ {C9E8E1E2-BFE3-41AD-8CB9-82F0F52E2306}.Release|Any CPU.Build.0 = Release|Any CPU
+ {C9E8E1E2-BFE3-41AD-8CB9-82F0F52E2306}.Release|x86.ActiveCfg = Release|Any CPU
+ {C9E8E1E2-BFE3-41AD-8CB9-82F0F52E2306}.Release|x86.Build.0 = Release|Any CPU
EndGlobalSection
GlobalSection(SolutionProperties) = preSolution
HideSolutionNode = FALSE