IFGame.Network Namespace - C# Network Connection Class
{
public class PacketInfo
{
private IWebSocket socket1;
public Header header;
public System.Object message = null;
public byte[] sendBuffer = null;
public byte[] messageBytes = null;
}
sealed public class Connection
{
#region Configuration
public delegate void CallbackD();
public delegate void OnConnectReport(bool success);
public CallbackD OnError;
public CallbackD OnClose;
public sealed class Configuration
{
public Configuration(string address,
int port,
int send_buffer_size = 1048576,
int recv_buffer_size = 1048576)
{
Address = address;
Port = port;
SendBufferSize = send_buffer_size;
ReceiveBufferSize = recv_buffer_size;
}
public string Address { get; private set; }
public int Port { get; private set; }
public int SendBufferSize { get; private set; }
public int ReceiveBufferSize { get; private set; }
}
#endregion
#region members
class SocketSession
{
public Socket socket;
public bool closed;
public SendAsyncState sendAsyncState;
public SocketSession(Socket socket)
{
this.socket = socket;
this.closed = false;
this.sendAsyncState = new SendAsyncState(this);
}
}
class ConnectAsyncState
{
public SocketSession socketSession;
public OnConnectReport report;
public ConnectAsyncState(SocketSession socketSession, OnConnectReport report)
{
this.socketSession = socketSession;
this.report = report;
}
}
class SendAsyncState
{
public SocketSession socketSession;
public int length;
public bool sending;
public SendAsyncState(SocketSession socketSession)
{
this.socketSession = socketSession;
this.length = 0;
this.sending = false;
}
}
class RecvAsyncState
{
public SocketSession socketSession;
public byte[] recvBuf;
public RecvAsyncState(SocketSession socketSession, int recvBufSize)
{
this.socketSession = socketSession;
this.recvBuf = new byte[recvBufSize];
}
}
Queue<PacketInfo> send_queue_ = new Queue<PacketInfo>();
Queue<PacketInfo> recv_queue_ = new Queue<PacketInfo>();
Codec codec_ = new Codec();
ChannelBufferLite channel_buffer_ = new ChannelBufferLite();
Configuration configuration_;
SocketSession socket_session_;
readonly object send_queue_lock_ = new object();
readonly object recv_queue_lock_ = new object();
#endregion
#region debug
// readonly object send_rec_lock_ = new object();
// LinkedList<float> send_rec_ = new LinkedList<float>();
int send_count_ = 0;
float clear_time_;
public float SendSpeed
{
get
{
float now = Time.time;
if (now - clear_time_ > 10)
{
send_count_ = 0;
clear_time_ = now;
return 0;
}
else
{
return send_count_ / (now - clear_time_);
}
}
}
#endregion debug
#region Encrypt
public void SetDefaultEncryptType(EncryptType type)
{
codec_.defaultEncrytType = type;
}
public void SetCommandEncryptType(uint command, EncryptType type)
{
if (!codec_.comandEncryptType.ContainsKey(command))
codec_.comandEncryptType.Add(command, type);
}
public void RemoveEncryptCommand(uint command)
{
codec_.comandEncryptType.Remove(command);
}
public void SetEncryptStaticKey(byte[] key)
{
codec_.EncryptStaticKey = key;
}
public void SetEncryptShareKey(byte[] key)
{
codec_.EncryptShareKey = key;
}
#endregion
#region public
public bool Connected { get { return null != socket_session_ ? socket_session_.socket.Connected : false; } }
public int RecvQueueCount { get { return recv_queue_.Count; } }
public Connection(Codec.Serialize serialize, Codec.Deserialize deserialize)
{
codec_.doSerialize = serialize;
codec_.doDeserialize = deserialize;
}
public void Close()
{
try
{
if (null != socket_session_)
{
socket_session_.closed = true;
socket_session_.socket.Shutdown(SocketShutdown.Both);
socket_session_.socket.Close();
socket_session_ = null;
}
}
catch (Exception e)
{
Debug.LogError('connection close exception:' + e.ToString());
}
}
public void Open(Configuration configuration, OnConnectReport connect_report)
{
try
{
Close();
//string addrees = '64:ff9b::cbc3:80e5';
Setup(configuration);
IPEndPoint endpoint = new IPEndPoint(
IPAddress.Parse(configuration_.Address),
configuration_.Port);
Debug.Log('connect to ' + configuration_.Address + ' : ' + configuration_.Port);
IAsyncResult result = socket_session_.socket.BeginConnect(
endpoint,
ConnectCallback,
new ConnectAsyncState(socket_session_, connect_report));
if (result == null)
throw new Exception('begin connect failed null result');
}
catch (Exception e)
{
connect_report(false);
Debug.LogError('connection open exception:' + e.ToString());
}
}
public void ClearAllMessage()
{
if (recv_queue_.Count > 0)
{
lock (recv_queue_lock_)
{
recv_queue_.Clear();
}
}
}
public PacketInfo NextReceived()
{
if (recv_queue_.Count > 0)
{
PacketInfo packet = null;
lock (recv_queue_lock_)
{
if (recv_queue_.Count > 0)
packet = recv_queue_.Dequeue();
}
return packet;
}
return null;
}
public void Send(PacketInfo packet, bool wait_replay = true)
{
try
{
if (packet.header.command > 0)
{
if (packet.messageBytes == null)
{
if (packet.message != null)
packet.messageBytes = codec_.Encode(packet.header, packet.message);
else if (packet.sendBuffer != null)
packet.messageBytes = codec_.Encode(packet.header, packet.sendBuffer);
}
}
lock (send_queue_lock_)
{
send_queue_.Enqueue(packet);
if (!socket_session_.sendAsyncState.sending)
PostAsyncSend(true, socket_session_);
}
}
catch (Exception e)
{
ErrorReport(e, 'connection send exception');
}
}
#endregion
#region private
void Setup(Configuration configuration)
{
configuration_ = configuration;
Socket socket = null;
Debug.LogWarning('socket address ' + configuration_.Address);
if(configuration_.Address.Contains(':'))
socket = new Socket(AddressFamily.InterNetworkV6, SocketType.Stream, ProtocolType.Tcp);
else
socket = new Socket(AddressFamily.InterNetwork, SocketType.Stream, ProtocolType.Tcp);
socket.SendBufferSize = configuration.SendBufferSize;
socket.ReceiveBufferSize = configuration.ReceiveBufferSize;
socket.NoDelay = true;
socket.LingerState = new LingerOption(true, 0);
socket_session_ = new SocketSession(socket);
send_queue_.Clear();
recv_queue_.Clear();
channel_buffer_.Clear();
}
void ConnectCallback(IAsyncResult ar)
{
ConnectAsyncState connAsyncState = (ConnectAsyncState)ar.AsyncState;
try
{
connAsyncState.socketSession.socket.EndConnect(ar);
connAsyncState.report(true);
RecvAsyncState recvAsyncState = new RecvAsyncState(
connAsyncState.socketSession,
configuration_.ReceiveBufferSize);
recvAsyncState.socketSession.socket.BeginReceive(
recvAsyncState.recvBuf,
0,
recvAsyncState.recvBuf.Length,
SocketFlags.None,
ReceiveCallback,
recvAsyncState);
}
catch (Exception e)
{
if (!connAsyncState.socketSession.closed)
{
connAsyncState.report(false);
Debug.LogWarning('connection ConnectCallback exception:' + e.ToString());
}
}
}
void ReceiveCallback(IAsyncResult ar)
{
RecvAsyncState recvAsyncState = (RecvAsyncState)ar.AsyncState;
try
{
int size = recvAsyncState.socketSession.socket.EndReceive(ar);
if (size == 0)
{
Debug.LogWarning('peer has closed the connection, we're disconnected');
if (socket_session_ == recvAsyncState.socketSession)
{
Close();
if (null != OnClose)
OnClose();
}
}
else
{
channel_buffer_.Write(recvAsyncState.recvBuf, size);
ExtractMessages();
recvAsyncState.socketSession.socket.BeginReceive(
recvAsyncState.recvBuf,
0,
recvAsyncState.recvBuf.Length,
SocketFlags.None,
ReceiveCallback, recvAsyncState);
}
}
catch (Exception e)
{
if (!recvAsyncState.socketSession.closed)
ErrorReport(e, 'connection ReceiveCallback exception');
}
}
void ExtractMessages()
{
while (true)
{
int offset;
int size;
byte[] buffer = channel_buffer_.GetBuffer(out offset, out size);
int bytes_read = 0;
try
{
Header header;
System.Object message;
byte[] message_bytes;
bool decode_ret = false;
decode_ret = codec_.Decode(buffer, offset, size, out bytes_read, out header, out message, out message_bytes);
if (!decode_ret)
break;
lock (recv_queue_lock_)
{
PacketInfo packet = new PacketInfo();
packet.header = header;
packet.message = message;
packet.messageBytes = message_bytes;
recv_queue_.Enqueue(packet);
}
}
catch (UnknownCommandException e)
{
Debug.LogException(e);
channel_buffer_.Consume(bytes_read);
break;
}
channel_buffer_.Consume(bytes_read);
}
}
bool PostAsyncSend(bool already_locked, SocketSession socket_session)
{
PacketInfo packet = null;
if (already_locked)
{
if (send_queue_.Count > 0)
packet = send_queue_.Dequeue();
}
else
{
lock (send_queue_lock_)
{
if (send_queue_.Count > 0)
packet = send_queue_.Dequeue();
}
}
if (packet != null)
{
socket_session_.sendAsyncState.sending = true;
if (packet.header.command == 0)
{
byte[] encoded_bytes = new byte[0];
socket_session_.sendAsyncState.length = encoded_bytes.Length;
socket_session_.socket.BeginSend(
encoded_bytes,
0,
encoded_bytes.Length,
SocketFlags.None,
SendCallback,
socket_session_.sendAsyncState);
}
else
{
socket_session_.sendAsyncState.length = packet.messageBytes.Length;
socket_session_.socket.BeginSend(
packet.messageBytes,
0,
packet.messageBytes.Length,
SocketFlags.None,
SendCallback,
socket_session_.sendAsyncState);
}
return true;
}
return false;
}
void SendCallback(IAsyncResult ar)
{
SendAsyncState sendAsyncState = (SendAsyncState)ar.AsyncState;
try
{
int sended_length = sendAsyncState.socketSession.socket.EndSend(ar);
if (sendAsyncState.length != sended_length)
{
throw new Exception('end send abnormal length ' + sendAsyncState.length.ToString() + ' ' + sended_length.ToString());
}
++send_count_;
if (!PostAsyncSend(false, sendAsyncState.socketSession))
sendAsyncState.sending = false;
}
catch (Exception e)
{
if (!sendAsyncState.socketSession.closed)
ErrorReport(e, 'connection SendCallback exception');
}
}
void ErrorReport(Exception e, string msg = null)
{
//for debug
if (msg != null && msg.Length != 0)
{
Debug.LogWarning(msg);
//IFGame.UI.UIManager.Instance.Notice(msg);
}
if (null != socket_session_ && !socket_session_.closed)
{
Debug.LogWarning(e.ToString());
if (null != OnError)
OnError();
}
}
#endregion
}
原文地址: https://www.cveoy.top/t/topic/qyI0 著作权归作者所有。请勿转载和采集!