{
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  
}
IFGame.Network Namespace - C# Network Connection Class

原文地址: https://www.cveoy.top/t/topic/qyI0 著作权归作者所有。请勿转载和采集!

免费AI点我,无需注册和登录