parallel delta queue processing

thread safe packet construction, compilation, compression, gump compilation, sending, coalescing, and many other fixes
TODO:  delta queue recursion fixups, ipooledenumerable fixups/generics
This commit is contained in:
Mark Sturgill 2013-10-11 00:01:05 -07:00
parent 9404651774
commit ddcc4e7a20
16 changed files with 619 additions and 575 deletions

View file

@ -83,92 +83,104 @@ namespace Server.Network {
// If our input exceeds this length, we may potentially overflow the buffer
private const int PossibleOverflow = ( ( BufferSize * 8 ) - TerminalCodeLength ) / MaximalCodeLength;
private static object _syncRoot = new object();
private static byte[] _outputBuffer = new byte[BufferSize];
[Obsolete( "Use Compress( byte[], int, int, ref int ) instead.", false )]
public static void Compress( byte[] input, int length, out byte[] output, out int outputLength ) {
outputLength = 0;
output = Compress( input, 0, length, ref outputLength );
}
public unsafe static byte[] Compress( byte[] input, int offset, int count, ref int length ) {
if ( input == null ) {
throw new ArgumentNullException( "input" );
} else if ( offset < 0 || offset >= input.Length ) {
throw new ArgumentOutOfRangeException( "offset" );
} else if ( count < 0 || count > input.Length ) {
throw new ArgumentOutOfRangeException( "count" );
} else if ( ( input.Length - offset ) < count ) {
public unsafe static void Compress(byte[] input, int offset, int count, byte[] output, ref int length)
{
if (input == null)
{
throw new ArgumentNullException("input");
}
else if (offset < 0 || offset >= input.Length)
{
throw new ArgumentOutOfRangeException("offset");
}
else if (count < 0 || count > input.Length)
{
throw new ArgumentOutOfRangeException("count");
}
else if ((input.Length - offset) < count)
{
throw new ArgumentException();
}
length = 0;
if ( count > DefiniteOverflow ) {
return null;
if (count > DefiniteOverflow)
{
return;
}
lock ( _syncRoot ) {
int bitCount = 0;
int bitValue = 0;
int bitCount = 0;
int bitValue = 0;
fixed ( int* pTable = _huffmanTable ) {
int* pEntry;
fixed (int* pTable = _huffmanTable)
{
int* pEntry;
fixed ( byte* pInputBuffer = input ) {
byte* pInput = pInputBuffer + offset, pInputEnd = pInput + count;
fixed (byte* pInputBuffer = input)
{
byte* pInput = pInputBuffer + offset, pInputEnd = pInput + count;
fixed ( byte* pOutputBuffer = _outputBuffer ) {
byte* pOutput = pOutputBuffer, pOutputEnd = pOutput + BufferSize;
fixed (byte* pOutputBuffer = output)
{
byte* pOutput = pOutputBuffer, pOutputEnd = pOutput + BufferSize;
while ( pInput < pInputEnd ) {
pEntry = &pTable[*pInput++ << 1];
bitCount += pEntry[CountIndex];
bitValue <<= pEntry[CountIndex];
bitValue |= pEntry[ValueIndex];
while ( bitCount >= 8 ) {
bitCount -= 8;
if ( pOutput < pOutputEnd ) {
*pOutput++ = ( byte ) ( bitValue >> bitCount );
} else {
return null;
}
}
}
// terminal code
pEntry = &pTable[0x200];
while (pInput < pInputEnd)
{
pEntry = &pTable[*pInput++ << 1];
bitCount += pEntry[CountIndex];
bitValue <<= pEntry[CountIndex];
bitValue |= pEntry[ValueIndex];
// align on byte boundary
if ( ( bitCount & 7 ) != 0 ) {
bitValue <<= ( 8 - ( bitCount & 7 ) );
bitCount += ( 8 - ( bitCount & 7 ) );
}
while ( bitCount >= 8 ) {
while (bitCount >= 8)
{
bitCount -= 8;
if ( pOutput < pOutputEnd ) {
*pOutput++ = ( byte ) ( bitValue >> bitCount );
} else {
return null;
if (pOutput < pOutputEnd)
{
*pOutput++ = (byte)(bitValue >> bitCount);
}
else
{
length = 0;
return;
}
}
length = ( int ) ( pOutput - pOutputBuffer );
return _outputBuffer;
}
// terminal code
pEntry = &pTable[0x200];
bitCount += pEntry[CountIndex];
bitValue <<= pEntry[CountIndex];
bitValue |= pEntry[ValueIndex];
// align on byte boundary
if ((bitCount & 7) != 0)
{
bitValue <<= (8 - (bitCount & 7));
bitCount += (8 - (bitCount & 7));
}
while (bitCount >= 8)
{
bitCount -= 8;
if (pOutput < pOutputEnd)
{
*pOutput++ = (byte)(bitValue >> bitCount);
}
else
{
length = 0;
return;
}
}
length = (int)(pOutput - pOutputBuffer);
return;
}
}
}

View file

@ -236,7 +236,9 @@ namespace Server.Network
return;
}
PacketReceiveProfile prof = PacketReceiveProfile.Acquire( packetID );
PacketReceiveProfile prof = null;
if (Core.Profiling) prof = PacketReceiveProfile.Acquire( packetID );
if ( prof != null ) {
prof.Start();

View file

@ -19,12 +19,14 @@
***************************************************************************/
using System;
using System.Collections;
using System.Collections.Generic;
using System.IO;
using System.Net;
using System.Net.Sockets;
using System.Threading;
#if Framework_4_0
using System.Threading.Tasks;
#endif
using Server;
using Server.Accounting;
using Server.Network;
@ -588,14 +590,15 @@ namespace Server.Network {
}
}
private bool _sending;
private object _sendL = new object();
public virtual void Send( Packet p ) {
if ( m_Socket == null || m_BlockAllPackets ) {
p.OnSend();
return;
}
PacketSendProfile prof = PacketSendProfile.Acquire( p.GetType() );
int length;
byte[] buffer = p.Compile( m_CompressionEnabled, out length );
@ -605,6 +608,10 @@ namespace Server.Network {
return;
}
PacketSendProfile prof = null;
if (Core.Profiling) prof = PacketSendProfile.Acquire(p.GetType());
if ( prof != null ) {
prof.Start();
}
@ -616,22 +623,27 @@ namespace Server.Network {
try {
SendQueue.Gram gram;
lock ( m_SendQueue ) {
gram = m_SendQueue.Enqueue( buffer, length );
}
lock (_sendL) {
lock (m_SendQueue)
gram = m_SendQueue.Enqueue(buffer, length);
if ( gram != null ) {
if (gram != null) {
#if NewAsyncSockets
m_SendEventArgs.SetBuffer( gram.Buffer, 0, gram.Length );
Send_Start();
m_SendEventArgs.SetBuffer( gram.Buffer, 0, gram.Length );
Send_Start();
#else
try {
m_Socket.BeginSend( gram.Buffer, 0, gram.Length, SocketFlags.None, m_OnSend, m_Socket );
} catch ( Exception ex ) {
TraceException( ex );
Dispose( false );
}
try {
if (!_sending) {
_sending = true;
m_Socket.BeginSend(gram.Buffer, 0, gram.Length, SocketFlags.None, m_OnSend, m_Socket);
}
}
catch (Exception ex) {
TraceException(ex);
Dispose(false);
}
#endif
}
}
} catch ( CapacityExceededException ) {
Console.WriteLine( "Client: {0}: Too much data pending, disconnecting...", this );
@ -814,13 +826,15 @@ namespace Server.Network {
}
public bool Flush() {
if ( m_Socket == null || !m_SendQueue.IsFlushReady ) {
return false;
}
if ( m_Socket == null )
return false;
SendQueue.Gram gram;
lock ( m_SendQueue ) {
if (!m_SendQueue.IsFlushReady)
return false;
gram = m_SendQueue.CheckFlushReady();
}
@ -914,23 +928,26 @@ namespace Server.Network {
m_NextCheckActivity = Core.TickCount + 90000;
if ( m_CoalesceSleep >= 0 ) {
Thread.Sleep( m_CoalesceSleep );
if (m_CoalesceSleep >= 0) {
Thread.Sleep(m_CoalesceSleep);
}
SendQueue.Gram gram;
lock ( m_SendQueue ) {
lock (m_SendQueue) {
gram = m_SendQueue.Dequeue();
}
if ( gram != null ) {
if (gram != null) {
try {
s.BeginSend( gram.Buffer, 0, gram.Length, SocketFlags.None, m_OnSend, s );
} catch ( Exception ex ) {
TraceException( ex );
Dispose( false );
s.BeginSend(gram.Buffer, 0, gram.Length, SocketFlags.None, m_OnSend, s);
} catch (Exception ex) {
TraceException(ex);
Dispose(false);
}
} else {
lock (_sendL)
_sending = false;
}
} catch ( Exception ){
Dispose( false );
@ -974,23 +991,31 @@ namespace Server.Network {
}
public bool Flush() {
if ( m_Socket == null || !m_SendQueue.IsFlushReady ) {
if (m_Socket == null)
return false;
}
SendQueue.Gram gram;
lock (_sendL) {
if (_sending)
return false;
lock ( m_SendQueue ) {
gram = m_SendQueue.CheckFlushReady();
}
SendQueue.Gram gram;
if ( gram != null ) {
try {
m_Socket.BeginSend( gram.Buffer, 0, gram.Length, SocketFlags.None, m_OnSend, m_Socket );
return true;
} catch ( Exception ex ) {
TraceException( ex );
Dispose( false );
lock (m_SendQueue) {
if (!m_SendQueue.IsFlushReady)
return false;
gram = m_SendQueue.CheckFlushReady();
}
if (gram != null) {
try {
_sending = true;
m_Socket.BeginSend(gram.Buffer, 0, gram.Length, SocketFlags.None, m_OnSend, m_Socket);
return true;
} catch (Exception ex) {
TraceException(ex);
Dispose(false);
}
}
}
@ -1007,11 +1032,13 @@ namespace Server.Network {
}
public static void FlushAll() {
#if Framework_4_0
Parallel.ForEach( m_Instances, ns => ns.Flush() );
#else
for ( int i = 0; i < m_Instances.Count; ++i ) {
NetState ns = m_Instances[i];
ns.Flush();
m_Instances[i].Flush();
}
#endif
}
private static int m_CoalesceSleep = -1;
@ -1109,10 +1136,11 @@ namespace Server.Network {
m_Running = false;
m_Disposed.Enqueue( this );
lock (m_Disposed)
m_Disposed.Enqueue( this );
if ( /*!flush &&*/ !m_SendQueue.IsEmpty ) {
lock ( m_SendQueue )
lock (m_SendQueue)
if ( /*!flush &&*/ !m_SendQueue.IsEmpty ) {
m_SendQueue.Clear();
}
}
@ -1131,37 +1159,38 @@ namespace Server.Network {
}
}
private static Queue m_Disposed = Queue.Synchronized( new Queue() );
private static Queue<NetState> m_Disposed = new Queue<NetState>();
public static void ProcessDisposedQueue() {
int breakout = 0;
lock (m_Disposed) {
int breakout = 0;
while ( breakout < 200 && m_Disposed.Count > 0 ) {
++breakout;
while ( breakout < 200 && m_Disposed.Count > 0 ) {
++breakout;
NetState ns = m_Disposed.Dequeue();
NetState ns = ( NetState ) m_Disposed.Dequeue();
Mobile m = ns.m_Mobile;
IAccount a = ns.m_Account;
Mobile m = ns.m_Mobile;
IAccount a = ns.m_Account;
if ( m != null ) {
m.NetState = null;
ns.m_Mobile = null;
}
if ( m != null ) {
m.NetState = null;
ns.m_Mobile = null;
}
ns.m_Gumps.Clear();
ns.m_Menus.Clear();
ns.m_HuePickers.Clear();
ns.m_Account = null;
ns.m_ServerInfo = null;
ns.m_CityInfo = null;
ns.m_Gumps.Clear();
ns.m_Menus.Clear();
ns.m_HuePickers.Clear();
ns.m_Account = null;
ns.m_ServerInfo = null;
ns.m_CityInfo = null;
m_Instances.Remove( ns );
m_Instances.Remove( ns );
if ( a != null ) {
ns.WriteConsole( "Disconnected. [{0} Online] [{1}]", m_Instances.Count, a );
} else {
ns.WriteConsole( "Disconnected. [{0} Online]", m_Instances.Count );
if ( a != null ) {
ns.WriteConsole( "Disconnected. [{0} Online] [{1}]", m_Instances.Count, a );
} else {
ns.WriteConsole( "Disconnected. [{0} Online]", m_Instances.Count );
}
}
}
}

View file

@ -97,7 +97,7 @@ namespace Server.Network
/// <summary>
/// Internal format buffer.
/// </summary>
private static byte[] m_Buffer = new byte[4];
private byte[] m_Buffer = new byte[4];
/// <summary>
/// Instantiates a new PacketWriter instance with the default capacity of 4 bytes.

View file

@ -2475,7 +2475,8 @@ namespace Server.Network
PacketWriter.ReleaseInstance( m_Strings );
}
private static byte[] m_PackBuffer;
private const int GumpBufferSize = 0x4000;
private static BufferPool m_PackBuffers = new BufferPool("Gump", 4, GumpBufferSize);
private void WritePacked( PacketWriter src )
{
@ -2493,8 +2494,15 @@ namespace Server.Network
wantLength += 4095;
wantLength &= ~4095;
if ( m_PackBuffer == null || m_PackBuffer.Length < wantLength )
byte[] m_PackBuffer;
lock (m_PackBuffers)
m_PackBuffer = m_PackBuffers.AcquireBuffer();
if (m_PackBuffer.Length < wantLength)
{
Console.WriteLine("Notice: DisplayGumpPacked creating new {0} byte buffer", wantLength);
m_PackBuffer = new byte[wantLength];
}
int packLength = m_PackBuffer.Length;
@ -2503,6 +2511,9 @@ namespace Server.Network
m_Stream.Write( (int) ( 4 + packLength ) );
m_Stream.Write( (int) length );
m_Stream.Write( m_PackBuffer, 0, packLength );
lock (m_PackBuffers)
m_PackBuffers.ReleaseBuffer(m_PackBuffer);
}
}
@ -2517,6 +2528,8 @@ namespace Server.Network
public DisplayGumpFast( Gump g ) : base( 0xB0 )
{
m_Buffer[0] = (byte)' ';
EnsureCapacity( 4096 );
m_Stream.Write( (int) g.Serial );
@ -2532,12 +2545,7 @@ namespace Server.Network
private static byte[] m_BeginTextSeparator = Gump.StringToBuffer( " @" );
private static byte[] m_EndTextSeparator = Gump.StringToBuffer( "@" );
private static byte[] m_Buffer = new byte[48];
static DisplayGumpFast()
{
m_Buffer[0] = (byte)' ';
}
private byte[] m_Buffer = new byte[48];
public void AppendLayout( bool val )
{
@ -4378,10 +4386,9 @@ namespace Server.Network
{
m_PacketID = packetID;
PacketSendProfile prof = PacketSendProfile.Acquire( GetType() );
if ( prof != null ) {
prof.Created++;
if (Core.Profiling) {
PacketSendProfile prof = PacketSendProfile.Acquire( GetType() );
prof.Increment();
}
}
@ -4400,10 +4407,9 @@ namespace Server.Network
m_Stream = PacketWriter.CreateInstance( length );// new PacketWriter( length );
m_Stream.Write( ( byte ) packetID );
PacketSendProfile prof = PacketSendProfile.Acquire( GetType() );
if ( prof != null ) {
prof.Created++;
if (Core.Profiling) {
PacketSendProfile prof = PacketSendProfile.Acquire( GetType() );
prof.Increment();
}
}
@ -4415,6 +4421,9 @@ namespace Server.Network
}
}
private const int CompressorBufferSize = 0x10000;
private static BufferPool m_CompressorBuffers = new BufferPool("Compressor", 4, CompressorBufferSize);
private const int BufferSize = 4096;
private static BufferPool m_Buffers = new BufferPool( "Compressed", 16, BufferSize );
@ -4490,8 +4499,10 @@ namespace Server.Network
{
Core.Set();
if ( (m_State & (State.Acquired | State.Static)) == 0 )
Free();
lock (this) {
if ((m_State & (State.Acquired | State.Static)) == 0)
Free();
}
}
private void Free()
@ -4499,8 +4510,8 @@ namespace Server.Network
if ( m_CompiledBuffer == null )
return;
if ( (m_State & State.Buffered) != 0 )
m_Buffers.ReleaseBuffer( m_CompiledBuffer );
if ((m_State & State.Buffered) != 0)
m_Buffers.ReleaseBuffer(m_CompiledBuffer);
m_State &= ~(State.Static | State.Acquired | State.Buffered);
@ -4509,7 +4520,7 @@ namespace Server.Network
public void Release()
{
if ( (m_State & State.Acquired) != 0 )
if ((m_State & State.Acquired) != 0)
Free();
}
@ -4518,43 +4529,46 @@ namespace Server.Network
public byte[] Compile( bool compress, out int length )
{
if ( m_CompiledBuffer == null )
lock (this)
{
if ( (m_State & State.Accessed) == 0 )
if (m_CompiledBuffer == null)
{
m_State |= State.Accessed;
}
else
{
if ( (m_State & State.Warned) == 0 )
if ((m_State & State.Accessed) == 0)
{
m_State |= State.Warned;
try
m_State |= State.Accessed;
}
else
{
if ((m_State & State.Warned) == 0)
{
using ( StreamWriter op = new StreamWriter( "net_opt.log", true ) )
m_State |= State.Warned;
try
{
using (StreamWriter op = new StreamWriter("net_opt.log", true))
{
op.WriteLine("Redundant compile for packet {0}, use Acquire() and Release()", this.GetType());
op.WriteLine(new System.Diagnostics.StackTrace());
}
}
catch
{
op.WriteLine( "Redundant compile for packet {0}, use Acquire() and Release()", this.GetType() );
op.WriteLine( new System.Diagnostics.StackTrace() );
}
}
catch
{
}
m_CompiledBuffer = new byte[0];
m_CompiledLength = 0;
length = m_CompiledLength;
return m_CompiledBuffer;
}
m_CompiledBuffer = new byte[0];
m_CompiledLength = 0;
length = m_CompiledLength;
return m_CompiledBuffer;
InternalCompile(compress);
}
InternalCompile( compress );
length = m_CompiledLength;
return m_CompiledBuffer;
}
length = m_CompiledLength;
return m_CompiledBuffer;
}
private void InternalCompile( bool compress )
@ -4578,41 +4592,49 @@ namespace Server.Network
m_CompiledBuffer = ms.GetBuffer();
int length = (int)ms.Length;
if ( compress )
{
m_CompiledBuffer = Compression.Compress(
m_CompiledBuffer, 0, length,
ref length
);
if ( m_CompiledBuffer == null )
{
Console.WriteLine( "Warning: Compression buffer overflowed on packet 0x{0:X2} ('{1}') (length={2})", m_PacketID, GetType().Name, length );
using ( StreamWriter op = new StreamWriter( "compression_overflow.log", true ) )
if ( compress ) {
byte[] buffer;
lock (m_CompressorBuffers)
buffer = m_CompressorBuffers.AcquireBuffer();
Compression.Compress(m_CompiledBuffer, 0, length, buffer, ref length);
if (length <= 0) {
Console.WriteLine("Warning: Compression buffer overflowed on packet 0x{0:X2} ('{1}') (length={2})", m_PacketID, GetType().Name, length);
using (StreamWriter op = new StreamWriter("compression_overflow.log", true))
{
op.WriteLine("{0} Warning: Compression buffer overflowed on packet 0x{1:X2} ('{2}') (length={3})", DateTime.UtcNow, m_PacketID, GetType().Name, length);
op.WriteLine( new System.Diagnostics.StackTrace() );
op.WriteLine(new System.Diagnostics.StackTrace());
}
}
}
} else {
m_CompiledLength = length;
if ( m_CompiledBuffer != null )
{
if (length > BufferSize || (m_State & State.Static) != 0) {
m_CompiledBuffer = new byte[length];
} else {
lock (m_Buffers)
m_CompiledBuffer = m_Buffers.AcquireBuffer();
m_State |= State.Buffered;
}
Buffer.BlockCopy(buffer, 0, m_CompiledBuffer, 0, length);
lock (m_CompressorBuffers)
m_CompressorBuffers.ReleaseBuffer(buffer);
}
} else if (length > 0) {
byte[] old = m_CompiledBuffer;
m_CompiledLength = length;
byte[] old = m_CompiledBuffer;
if ( length > BufferSize || (m_State & State.Static) != 0 )
{
if ( length > BufferSize || (m_State & State.Static) != 0 ) {
m_CompiledBuffer = new byte[length];
}
else
{
m_CompiledBuffer = m_Buffers.AcquireBuffer();
} else {
lock (m_Buffers)
m_CompiledBuffer = m_Buffers.AcquireBuffer();
m_State |= State.Buffered;
}
Buffer.BlockCopy( old, 0, m_CompiledBuffer, 0, length );
Buffer.BlockCopy(old, 0, m_CompiledBuffer, 0, length);
}
PacketWriter.ReleaseInstance( m_Stream );

View file

@ -104,22 +104,27 @@ namespace Server.Network {
if ( m_CoalesceBufferSize == value )
return;
if ( m_UnusedBuffers != null )
m_UnusedBuffers.Free();
BufferPool old = m_UnusedBuffers;
m_CoalesceBufferSize = value;
m_UnusedBuffers = new BufferPool( "Coalesced", 2048, m_CoalesceBufferSize );
lock (old) {
if ( m_UnusedBuffers != null )
m_UnusedBuffers.Free();
m_CoalesceBufferSize = value;
m_UnusedBuffers = new BufferPool( "Coalesced", 2048, m_CoalesceBufferSize );
}
}
}
public static byte[] AcquireBuffer() {
return m_UnusedBuffers.AcquireBuffer();
lock (m_UnusedBuffers)
return m_UnusedBuffers.AcquireBuffer();
}
public static void ReleaseBuffer( byte[] buffer ) {
if ( buffer != null && buffer.Length == m_CoalesceBufferSize ) {
m_UnusedBuffers.ReleaseBuffer( buffer );
}
lock (m_UnusedBuffers)
if ( buffer != null && buffer.Length == m_CoalesceBufferSize )
m_UnusedBuffers.ReleaseBuffer( buffer );
}
private Queue<Gram> _pending;
@ -143,15 +148,9 @@ namespace Server.Network {
}
public Gram CheckFlushReady() {
Gram gram = null;
if ( _pending.Count == 0 && _buffered != null ) {
gram = _buffered;
_pending.Enqueue( _buffered );
_buffered = null;
}
Gram gram = _buffered;
_pending.Enqueue(_buffered);
_buffered = null;
return gram;
}
@ -169,7 +168,7 @@ namespace Server.Network {
return gram;
}
private const int PendingCap = 96 * 1024;
private const int PendingCap = 128 * 1024;
public Gram Enqueue( byte[] buffer, int length ) {
return Enqueue( buffer, 0, length );