Switches to Kestrel (#102)

This commit is contained in:
Kamron Batman 2020-04-12 20:57:48 -07:00 • committed by GitHub
parent 70fddfebde
commit 2c0d97cd36
No known key found for this signature in database
GPG key ID: 4AEE18F83AFDEB23
141 changed files with 1040 additions and 1679 deletions

View file

@ -1,177 +0,0 @@
/***************************************************************************
* Listener.cs
* -------------------
* begin : May 1, 2002
* copyright : (C) The RunUO Software Team
* email : info@runuo.com
*
* $Id$
*
***************************************************************************/
/***************************************************************************
*
* This program is free software; you can redistribute it and/or modify
* it under the terms of the GNU General Public License as published by
* the Free Software Foundation; either version 2 of the License, or
* (at your option) any later version.
*
***************************************************************************/
using System;
using System.Net;
using System.Net.NetworkInformation;
using System.Net.Sockets;
using System.Threading;
using System.Threading.Tasks;
using Libuv;
using Libuv.Internal;
using Microsoft.AspNetCore.Connections;
using Microsoft.AspNetCore.Hosting;
using Microsoft.Extensions.Logging;
namespace Server.Network
{
public class Listener
{
private static readonly LibuvFunctions functions = new LibuvFunctions();
private readonly IPEndPoint m_EndPoint;
private LibuvConnectionListener m_Listener;
public Listener(IPEndPoint ipep)
{
m_EndPoint = ipep;
LibuvTransportContext transport = new LibuvTransportContext
{
Options = new LibuvTransportOptions(),
AppLifetime = new ApplicationLifetime(
LoggerFactory.Create(builder => { builder.AddConsole(); }).CreateLogger<ApplicationLifetime>()
),
Log = new LibuvTrace(LoggerFactory.Create(builder => { builder.AddConsole(); }).CreateLogger("network"))
};
m_Listener = new LibuvConnectionListener(functions, transport, ipep);
}
public virtual async Task Start(MessagePump pump)
{
try
{
await m_Listener.BindAsync();
}
catch (AddressInUseException)
{
Console.WriteLine("Listener Failed: {0}:{1} (In Use)", m_EndPoint.Address, m_EndPoint.Port);
m_Listener = null;
return;
}
catch (Exception e)
{
Console.WriteLine("Listener Exception:");
Console.WriteLine(e);
m_Listener = null;
return;
}
DisplayListener();
while (true)
{
ConnectionContext context;
try
{
context = await m_Listener.AcceptAsync();
}
catch (SocketException ex)
{
NetState.TraceException(ex);
continue;
}
if (VerifySocket(context))
_ = new NetState(context, pump);
else
Release(context);
}
}
private void DisplayListener()
{
if (!(m_Listener.EndPoint is IPEndPoint ipep))
return;
if (ipep.Address.Equals(IPAddress.Any) || ipep.Address.Equals(IPAddress.IPv6Any))
{
NetworkInterface[] adapters = NetworkInterface.GetAllNetworkInterfaces();
foreach (NetworkInterface adapter in adapters)
{
IPInterfaceProperties properties = adapter.GetIPProperties();
foreach (UnicastIPAddressInformation unicast in properties.UnicastAddresses)
if (ipep.AddressFamily == unicast.Address.AddressFamily)
Console.WriteLine("Listening: {0}:{1}", unicast.Address, ipep.Port);
}
}
else
Console.WriteLine("Listening: {0}:{1}", ipep.Address, ipep.Port);
}
private static bool VerifySocket(ConnectionContext context)
{
try
{
SocketConnectEventArgs args = new SocketConnectEventArgs(context);
EventSink.InvokeSocketConnect(args);
return args.AllowConnection;
}
catch (Exception ex)
{
NetState.TraceException(ex);
return false;
}
}
private static void Release(ConnectionContext context)
{
try
{
context.Abort(new ConnectionAbortedException("Failed socket verification."));
}
catch (Exception ex)
{
NetState.TraceException(ex);
}
try
{
// TODO: Is this needed?
context.DisposeAsync();
}
catch (Exception ex)
{
NetState.TraceException(ex);
}
}
public async Task Dispose()
{
LibuvConnectionListener listener = Interlocked.Exchange(ref m_Listener, null);
if (listener != null)
try
{
await listener.UnbindAsync();
await listener.DisposeAsync();
}
catch (Exception ex)
{
Console.WriteLine("Listener: Failed to dispose.");
Console.WriteLine(ex);
}
GC.SuppressFinalize(this);
}
}
}

View file

@ -19,26 +19,20 @@
*
***************************************************************************/
using System;
using System.Buffers;
using System.Collections.Concurrent;
using System.Net;
namespace Server.Network
{
public class MessagePump
public interface IMessagePumpService
{
private ConcurrentQueue<Work> m_WorkQueue = new ConcurrentQueue<Work>();
public Listener[] Listeners => new Listener[0];
void QueueWork(NetState ns, IMemoryOwner<byte> memOwner, OnPacketReceive onReceive);
void DoWork();
}
public void AddListener(IPEndPoint ipep)
{
Listener[] listeners = new Listener[Listeners.Length + 1];
Array.Copy(Listeners, listeners, Listeners.Length);
Listener listener = new Listener(ipep);
_ = listener.Start(this);
listeners[Listeners.Length] = listener;
}
public class MessagePumpService : IMessagePumpService
{
private readonly ConcurrentQueue<Work> m_WorkQueue = new ConcurrentQueue<Work>();
public void QueueWork(NetState ns, IMemoryOwner<byte> memOwner, OnPacketReceive onReceive)
{

View file

@ -19,7 +19,6 @@
***************************************************************************/
using System;
using System.Buffers;
using System.Collections.Concurrent;
using System.Collections.Generic;
using System.IO;
@ -91,7 +90,7 @@ namespace Server.Network
public class NetState : IComparable<NetState>
{
private string m_ToString;
private readonly string m_ToString;
private ClientVersion m_Version;
public DateTime ConnectedOn { get; }
@ -105,10 +104,11 @@ namespace Server.Network
public IPAddress Address { get; }
private static AsyncState m_PauseState = new AsyncState(true);
private static AsyncState m_ResumeState = new AsyncState(false);
private static readonly AsyncState m_PauseState = new AsyncState(true);
private static readonly AsyncState m_ResumeState = new AsyncState(false);
private static AsyncState m_AsyncState = m_ResumeState;
public static AsyncState AsyncState => m_AsyncState;
public IPacketEncoder PacketEncoder { get; set; }
@ -279,9 +279,6 @@ namespace Server.Network
public List<IMenu> Menus { get; private set; }
public PipeReader RecvPipe => Connection.Transport.Input;
public PipeWriter SendPipe => Connection.Transport.Output;
public static int GumpCap { get; set; } = 512;
public static int HuePickerCap { get; set; } = 512;
@ -399,11 +396,7 @@ namespace Server.Network
public override string ToString() => m_ToString;
public static List<NetState> Instances { get; } = new List<NetState>();
public void SetConnectionAlive() => m_NextCheckActivity = Core.TickCount + 60000;
public NetState(ConnectionContext connection, MessagePump pump)
public NetState(ConnectionContext connection)
{
Connection = connection;
Seeded = false;
@ -412,10 +405,6 @@ namespace Server.Network
Menus = new List<IMenu>();
Trades = new List<SecureTrade>();
SetConnectionAlive();
Instances.Add(this);
try
{
Address = Utility.Intern(((IPEndPoint)Connection.RemoteEndPoint).Address);
@ -429,11 +418,6 @@ namespace Server.Network
}
ConnectedOn = DateTime.UtcNow;
Console.WriteLine("Client: {0}: Connected. [{1} Online]", this, Instances.Count);
#pragma warning disable 4014
ProcessRecvs(pump);
#pragma warning restore 4014
CreatedCallback?.Invoke(this);
}
@ -448,7 +432,7 @@ namespace Server.Network
m_AsyncState = Interlocked.Exchange(ref m_AsyncState, m_ResumeState);
}
public virtual void Send(Packet p)
public virtual async void Send(Packet p)
{
if (Connection == null || BlockAllPackets)
{
@ -456,19 +440,29 @@ namespace Server.Network
return;
}
var outPipe = Connection.Transport.Output;
try
{
ReadOnlySpan<byte> buffer = p.Compile(CompressionEnabled, out int length);
//TODO: Rented memory
ReadOnlyMemory<byte> buffer = p.Compile(CompressionEnabled, out int length);
if (buffer.Length <= 0 || length <= 0)
{
p.OnSend();
return;
}
if (buffer.Length > 0 && length > 0)
try
{
FlushResult result = await outPipe.WriteAsync(buffer.Slice(0, length));
buffer.Slice(0, length).CopyTo(SendPipe.GetSpan(length));
SendPipe.Advance(length);
SendPipe.FlushAsync().GetAwaiter().GetResult();
if (result.IsCanceled || result.IsCompleted)
{
Dispose();
return;
}
}
catch
{
Dispose();
// ignored
}
p.OnSend();
}
@ -497,48 +491,10 @@ namespace Server.Network
return false;
}
private async Task ProcessRecvs(MessagePump pump)
{
while (true)
{
ReadResult result = await RecvPipe.ReadAsync();
ReadOnlySequence<byte> seq = result.Buffer;
if (seq.IsEmpty)
break;
SetConnectionAlive();
int pos = PacketHandlers.ProcessPacket(pump, this, seq);
if (pos <= 0)
break;
RecvPipe.AdvanceTo(seq.Slice(0, pos).End);
if (result.IsCompleted || result.IsCanceled)
break;
}
RecvPipe.Complete();
Dispose();
}
public PacketHandler GetHandler(int packetID) =>
ContainerGridLines ? PacketHandlers.Get6017Handler(packetID) :
PacketHandlers.GetHandler(packetID);
private long m_NextCheckActivity;
public void CheckAlive(long curTicks)
{
if (Connection == null || m_NextCheckActivity - curTicks >= 0)
return;
Console.WriteLine("Client: {0}: Disconnecting due to inactivity...", this);
Dispose();
}
public static void TraceException(Exception ex)
{
if (!Core.Debug)
@ -586,30 +542,7 @@ namespace Server.Network
m_Disposed.Enqueue(this);
}
public static void Initialize()
{
Timer.DelayCall(TimeSpan.FromMinutes(1.0), TimeSpan.FromMinutes(1.5), CheckAllAlive);
}
public static void CheckAllAlive()
{
try
{
long curTicks = Core.TickCount;
if (Instances.Count >= 1024)
Parallel.ForEach(Instances, ns => ns.CheckAlive(curTicks));
else
for (int i = 0; i < Instances.Count; ++i)
Instances[i].CheckAlive(curTicks);
}
catch (Exception ex)
{
TraceException(ex);
}
}
private static ConcurrentQueue<NetState> m_Disposed = new ConcurrentQueue<NetState>();
private static readonly ConcurrentQueue<NetState> m_Disposed = new ConcurrentQueue<NetState>();
public static void ProcessDisposedQueue()
{
@ -636,12 +569,10 @@ namespace Server.Network
ns.ServerInfo = null;
ns.CityInfo = null;
Instances.Remove(ns);
if (a != null)
ns.WriteConsole("Disconnected. [{0} Online] [{1}]", Instances.Count, a);
ns.WriteConsole("Disconnected. [{0} Online] [{1}]", TcpServer.Instances.Count, a);
else
ns.WriteConsole("Disconnected. [{0} Online]", Instances.Count);
ns.WriteConsole("Disconnected. [{0} Online]", TcpServer.Instances.Count);
}
}

View file

@ -34,7 +34,7 @@ namespace Server.Network
private byte[] m_CompiledBuffer;
private int m_CompiledLength;
private int m_Length;
private readonly int m_Length;
private State m_State;
protected PacketWriter m_Stream;

View file

@ -62,17 +62,17 @@ namespace Server.Network
private const int BadUOTD = unchecked((int)0xFFCEFFCE);
private const int m_AuthIDWindowSize = 128;
private static PacketHandler[] m_6017Handlers;
private static readonly PacketHandler[] m_6017Handlers;
private static PacketHandler[] m_ExtendedHandlersLow;
private static Dictionary<int, PacketHandler> m_ExtendedHandlersHigh;
private static readonly PacketHandler[] m_ExtendedHandlersLow;
private static readonly Dictionary<int, PacketHandler> m_ExtendedHandlersHigh;
private static EncodedPacketHandler[] m_EncodedHandlersLow;
private static Dictionary<int, EncodedPacketHandler> m_EncodedHandlersHigh;
private static readonly EncodedPacketHandler[] m_EncodedHandlersLow;
private static readonly Dictionary<int, EncodedPacketHandler> m_EncodedHandlersHigh;
private static int[] m_EmptyInts = new int[0];
private static readonly int[] m_EmptyInts = new int[0];
private static KeywordList m_KeywordList = new KeywordList();
private static readonly KeywordList m_KeywordList = new KeywordList();
public static int[] m_ValidAnimations =
{
@ -91,7 +91,7 @@ namespace Server.Network
public static PlayCharCallback ThirdPartyAuthCallback = null, ThirdPartyHackedCallback = null;
private static Dictionary<int, AuthIDPersistence> m_AuthIDWindow =
private static readonly Dictionary<int, AuthIDPersistence> m_AuthIDWindow =
new Dictionary<int, AuthIDPersistence>(m_AuthIDWindowSize);
static PacketHandlers()
@ -277,9 +277,9 @@ namespace Server.Network
ph.ThrottleCallback = t;
}
private static MemoryPool<byte> _memoryPool = SlabMemoryPoolFactory.Create();
private static readonly MemoryPool<byte> _memoryPool = SlabMemoryPoolFactory.Create();
public static int ProcessPacket(MessagePump pump, NetState ns, in ReadOnlySequence<byte> seq)
public static int ProcessPacket(IMessagePumpService pump, NetState ns, in ReadOnlySequence<byte> seq)
{
PacketReader r = new PacketReader(seq);
@ -2495,8 +2495,8 @@ namespace Server.Network
private class LoginTimer : Timer
{
private Mobile m_Mobile;
private NetState m_State;
private readonly Mobile m_Mobile;
private readonly NetState m_State;
public LoginTimer(NetState state, Mobile m) : base(TimeSpan.FromSeconds(1.0), TimeSpan.FromSeconds(1.0))
{
@ -2507,7 +2507,10 @@ namespace Server.Network
protected override void OnTick()
{
if (m_State == null)
{
Stop();
return;
}
if (m_State.Version != null)
{

View file

@ -30,12 +30,12 @@ namespace Server.Network
/// </summary>
public class PacketWriter
{
private static ConcurrentQueue<PacketWriter> m_Pool = new ConcurrentQueue<PacketWriter>();
private static readonly ConcurrentQueue<PacketWriter> m_Pool = new ConcurrentQueue<PacketWriter>();
/// <summary>
/// Internal format buffer.
/// </summary>
private byte[] m_Buffer = new byte[4];
private readonly byte[] m_Buffer = new byte[4];
private int m_Capacity;

View file

@ -540,7 +540,7 @@ namespace Server.Network
public sealed class ChangeUpdateRange : Packet
{
private static ChangeUpdateRange[] m_Cache = new ChangeUpdateRange[0x100];
private static readonly ChangeUpdateRange[] m_Cache = new ChangeUpdateRange[0x100];
public ChangeUpdateRange(int range) : base(0xC8, 2)
{
@ -811,7 +811,7 @@ namespace Server.Network
public sealed class GlobalLightLevel : Packet
{
private static GlobalLightLevel[] m_Cache = new GlobalLightLevel[0x100];
private static readonly GlobalLightLevel[] m_Cache = new GlobalLightLevel[0x100];
public GlobalLightLevel(int level) : base(0x4F, 2)
{
@ -2105,9 +2105,9 @@ namespace Server.Network
public sealed class MessageLocalized : Packet
{
private static MessageLocalized[] m_Cache_IntLoc = new MessageLocalized[15000];
private static MessageLocalized[] m_Cache_CliLoc = new MessageLocalized[100000];
private static MessageLocalized[] m_Cache_CliLocCmp = new MessageLocalized[5000];
private static readonly MessageLocalized[] m_Cache_IntLoc = new MessageLocalized[15000];
private static readonly MessageLocalized[] m_Cache_CliLoc = new MessageLocalized[100000];
private static readonly MessageLocalized[] m_Cache_CliLocCmp = new MessageLocalized[5000];
public MessageLocalized(Serial serial, int graphic, MessageType type, int hue, int font, int number, string name,
string args) : base(0xC1)
@ -2318,20 +2318,20 @@ namespace Server.Network
public sealed class DisplayGumpPacked : Packet, IGumpWriter
{
private static byte[] m_True = Gump.StringToBuffer(" 1");
private static byte[] m_False = Gump.StringToBuffer(" 0");
private static readonly byte[] m_True = Gump.StringToBuffer(" 1");
private static readonly byte[] m_False = Gump.StringToBuffer(" 0");
private static byte[] m_BeginTextSeparator = Gump.StringToBuffer(" @");
private static byte[] m_EndTextSeparator = Gump.StringToBuffer("@");
private static readonly byte[] m_BeginTextSeparator = Gump.StringToBuffer(" @");
private static readonly byte[] m_EndTextSeparator = Gump.StringToBuffer("@");
private static byte[] m_Buffer = new byte[48];
private static readonly byte[] m_Buffer = new byte[48];
private Gump m_Gump;
private readonly Gump m_Gump;
private PacketWriter m_Layout;
private readonly PacketWriter m_Layout;
private int m_StringCount;
private PacketWriter m_Strings;
private readonly PacketWriter m_Strings;
static DisplayGumpPacked() => m_Buffer[0] = (byte)' ';
@ -2457,13 +2457,13 @@ namespace Server.Network
public sealed class DisplayGumpFast : Packet, IGumpWriter
{
private static byte[] m_True = Gump.StringToBuffer(" 1");
private static byte[] m_False = Gump.StringToBuffer(" 0");
private static readonly byte[] m_True = Gump.StringToBuffer(" 1");
private static readonly byte[] m_False = Gump.StringToBuffer(" 0");
private static byte[] m_BeginTextSeparator = Gump.StringToBuffer(" @");
private static byte[] m_EndTextSeparator = Gump.StringToBuffer("@");
private static readonly byte[] m_BeginTextSeparator = Gump.StringToBuffer(" @");
private static readonly byte[] m_EndTextSeparator = Gump.StringToBuffer("@");
private byte[] m_Buffer = new byte[48];
private readonly byte[] m_Buffer = new byte[48];
private int m_LayoutLength;
public DisplayGumpFast(Gump g) : base(0xB0)
@ -2629,7 +2629,7 @@ namespace Server.Network
{
public static readonly Packet InvalidInstance = SetStatic(new PlayMusic(MusicName.Invalid));
private static Packet[] m_Instances = new Packet[60];
private static readonly Packet[] m_Instances = new Packet[60];
public PlayMusic(MusicName name) : base(0x6D, 3)
{
@ -2700,7 +2700,7 @@ namespace Server.Network
public sealed class SeasonChange : Packet
{
private static SeasonChange[][] m_Cache = new SeasonChange[][]
private static readonly SeasonChange[][] m_Cache = new SeasonChange[][]
{
new SeasonChange[2],
new SeasonChange[2],
@ -3250,8 +3250,8 @@ namespace Server.Network
public sealed class MobileIncoming : Packet
{
private static ThreadLocal<int[]> m_DupedLayersTL = new ThreadLocal<int[]>(() => { return new int[256]; });
private static ThreadLocal<int> m_VersionTL = new ThreadLocal<int>();
private static readonly ThreadLocal<int[]> m_DupedLayersTL = new ThreadLocal<int[]>(() => { return new int[256]; });
private static readonly ThreadLocal<int> m_VersionTL = new ThreadLocal<int>();
public Mobile m_Beheld;
@ -3363,8 +3363,8 @@ namespace Server.Network
public sealed class MobileIncomingSA : Packet
{
private static ThreadLocal<int[]> m_DupedLayersTL = new ThreadLocal<int[]>(() => { return new int[256]; });
private static ThreadLocal<int> m_VersionTL = new ThreadLocal<int>();
private static readonly ThreadLocal<int[]> m_DupedLayersTL = new ThreadLocal<int[]>(() => { return new int[256]; });
private static readonly ThreadLocal<int> m_VersionTL = new ThreadLocal<int>();
public Mobile m_Beheld;
@ -3485,8 +3485,8 @@ namespace Server.Network
// Pre-7.0.0.0 Mobile Incoming
public sealed class MobileIncomingOld : Packet
{
private static ThreadLocal<int[]> m_DupedLayersTL = new ThreadLocal<int[]>(() => { return new int[256]; });
private static ThreadLocal<int> m_VersionTL = new ThreadLocal<int>();
private static readonly ThreadLocal<int[]> m_DupedLayersTL = new ThreadLocal<int[]>(() => { return new int[256]; });
private static readonly ThreadLocal<int> m_VersionTL = new ThreadLocal<int>();
public Mobile m_Beheld;
@ -3654,7 +3654,7 @@ namespace Server.Network
public sealed class PingAck : Packet
{
private static PingAck[] m_Cache = new PingAck[0x100];
private static readonly PingAck[] m_Cache = new PingAck[0x100];
public PingAck(byte ping) : base(0x73, 2)
{
@ -3689,7 +3689,7 @@ namespace Server.Network
public sealed class MovementAck : Packet
{
private static MovementAck[] m_Cache = new MovementAck[8 * 256];
private static readonly MovementAck[] m_Cache = new MovementAck[8 * 256];
private MovementAck(int seq, int noto) : base(0x22, 3)
{

View file

@ -0,0 +1,23 @@
using System.Threading.Tasks;
using Microsoft.AspNetCore.Builder;
using Microsoft.Extensions.DependencyInjection;
using Server.Network;
namespace Server.Network
{
public class ServerStartup
{
private readonly IMessagePumpService _messagePumpService;
public ServerStartup(IMessagePumpService messagePumpService) => _messagePumpService = messagePumpService;
public void ConfigureServices(IServiceCollection services)
{
}
public void Configure(IApplicationBuilder app)
{
// Run async?
Task.Run(() => Core.RunEventLoop(_messagePumpService));
}
}
}

View file

@ -19,19 +19,18 @@
* along with this program. If not, see <http://www.gnu.org/licenses/>. *
*************************************************************************/
using System;
using System.Collections.Concurrent;
namespace Server.Network
{
public static class StaticPacketHandlers
{
private static ConcurrentDictionary<IPropertyListObject,OPLInfo> OPLInfoPackets = new ConcurrentDictionary<IPropertyListObject,OPLInfo>();
private static ConcurrentDictionary<IEntity,RemoveEntity> RemoveEntityPackets = new ConcurrentDictionary<IEntity,RemoveEntity>();
private static readonly ConcurrentDictionary<IPropertyListObject,OPLInfo> OPLInfoPackets = new ConcurrentDictionary<IPropertyListObject,OPLInfo>();
private static readonly ConcurrentDictionary<IEntity,RemoveEntity> RemoveEntityPackets = new ConcurrentDictionary<IEntity,RemoveEntity>();
private static ConcurrentDictionary<Item,WorldItem> WorldItemPackets = new ConcurrentDictionary<Item,WorldItem>();
private static ConcurrentDictionary<Item,WorldItemSA> WorldItemSAPackets = new ConcurrentDictionary<Item,WorldItemSA>();
private static ConcurrentDictionary<Item,WorldItemHS> WorldItemHSPackets = new ConcurrentDictionary<Item,WorldItemHS>();
private static readonly ConcurrentDictionary<Item,WorldItem> WorldItemPackets = new ConcurrentDictionary<Item,WorldItem>();
private static readonly ConcurrentDictionary<Item,WorldItemSA> WorldItemSAPackets = new ConcurrentDictionary<Item,WorldItemSA>();
private static readonly ConcurrentDictionary<Item,WorldItemHS> WorldItemHSPackets = new ConcurrentDictionary<Item,WorldItemHS>();
public static OPLInfo GetOPLInfoPacket(IPropertyListObject obj)
{

View file

@ -0,0 +1,57 @@
using System;
using System.Collections.Generic;
using System.Net;
using System.Net.NetworkInformation;
using Microsoft.AspNetCore;
using Microsoft.AspNetCore.Connections;
using Microsoft.AspNetCore.Hosting;
using Microsoft.Extensions.DependencyInjection;
namespace Server.Network
{
public class TcpServer
{
public static List<IPEndPoint> Listeners { get; } = new List<IPEndPoint>();
// Make this thread safe
public static List<NetState> Instances { get; } = new List<NetState>();
public static IWebHostBuilder CreateWebHostBuilder(string[] args) =>
WebHost.CreateDefaultBuilder(args)
.UseSetting(WebHostDefaults.SuppressStatusMessagesKey, "True")
.ConfigureServices(services =>
{
services.AddSingleton<IMessagePumpService>(new MessagePumpService());
})
.UseKestrel(options =>
{
foreach (var ipep in Listeners)
{
options.Listen(ipep, builder => { builder.UseConnectionHandler<ServerConnectionHandler>(); });
DisplayListener(ipep);
}
options.ListenLocalhost(2593, builder => { builder.UseConnectionHandler<ServerConnectionHandler>(); });
// Webservices here
})
.UseLibuv()
.UseStartup<ServerStartup>();
private static void DisplayListener(IPEndPoint ipep)
{
if (ipep.Address.Equals(IPAddress.Any) || ipep.Address.Equals(IPAddress.IPv6Any))
{
NetworkInterface[] adapters = NetworkInterface.GetAllNetworkInterfaces();
foreach (NetworkInterface adapter in adapters)
{
IPInterfaceProperties properties = adapter.GetIPProperties();
foreach (UnicastIPAddressInformation unicast in properties.UnicastAddresses)
if (ipep.AddressFamily == unicast.Address.AddressFamily)
Console.WriteLine("Listening: {0}:{1}", unicast.Address, ipep.Port);
}
}
else
Console.WriteLine("Listening: {0}:{1}", ipep.Address, ipep.Port);
}
}
}

View file

@ -0,0 +1,116 @@
using System;
using System.Buffers;
using System.IO.Pipelines;
using System.Threading.Tasks;
using Microsoft.AspNetCore.Connections;
using Microsoft.Extensions.Logging;
namespace Server.Network
{
public class ServerConnectionHandler : ConnectionHandler
{
private readonly IMessagePumpService _messagePumpService;
private readonly ILogger<ServerConnectionHandler> _logger;
public ServerConnectionHandler(
IMessagePumpService messagePumpService,
ILogger<ServerConnectionHandler> logger
)
{
_messagePumpService = messagePumpService;
_logger = logger;
}
public override async Task OnConnectedAsync(ConnectionContext connection)
{
if (!VerifySocket(connection))
{
Release(connection);
return;
}
NetState ns = new NetState(connection);
TcpServer.Instances.Add(ns);
_logger.LogInformation($"Client: {ns}: Connected. [{TcpServer.Instances.Count} Online]");
connection.ConnectionClosed.Register(() => { TcpServer.Instances.Remove(ns); });
await ProcessIncoming(ns);
}
private async Task ProcessIncoming(NetState ns)
{
var inPipe = ns.Connection.Transport.Input;
while (true)
{
if (NetState.AsyncState.Paused)
continue;
try
{
ReadResult result = await inPipe.ReadAsync();
if (result.IsCanceled || result.IsCompleted)
return;
ReadOnlySequence<byte> seq = result.Buffer;
if (seq.IsEmpty)
break;
int pos = PacketHandlers.ProcessPacket(_messagePumpService, ns, seq);
if (pos <= 0)
break;
inPipe.AdvanceTo(seq.Slice(0, pos).End);
}
catch
{
// ignored
}
}
inPipe.Complete();
}
private static bool VerifySocket(ConnectionContext connection)
{
try
{
SocketConnectEventArgs args = new SocketConnectEventArgs(connection);
EventSink.InvokeSocketConnect(args);
return args.AllowConnection;
}
catch (Exception ex)
{
NetState.TraceException(ex);
return false;
}
}
private static void Release(ConnectionContext connection)
{
try
{
connection.Abort(new ConnectionAbortedException("Failed socket verification."));
}
catch (Exception ex)
{
NetState.TraceException(ex);
}
try
{
// TODO: Is this needed?
connection.DisposeAsync();
}
catch (Exception ex)
{
NetState.TraceException(ex);
}
}
}
}