【问题标题】:Waiting on Interlocked == 0?等待联锁 == 0?
【发布时间】:2018-04-28 22:09:29
【问题描述】:

免责声明:我的 C# 甚至不如我的 C++

我正在尝试学习如何在 C# 中执行异步套接字,以便为我的组件编写测试应用程序。我以前使用 TcpClient 的尝试以失败告终,您可以在此处阅读未解决的问题:

TcpClient.NetworkStream Async operations - Canceling / Disconnect

Detect errors with NetworkStream.WriteAsync

因为我无法让它工作,所以我尝试改用 Socket.BeginX 和 Socket.EndX。我走得更远了。我现在的问题是,在下面的清单中,当需要断开连接时,这又会在套接字上调用 shutdown 和 close,异步操作仍然未完成,它们将抛出对象已处置异常或对象设置为空异常。

我在这里找到了一个类似的帖子:

After disposing async socket (.Net) callbacks still get called

但是,我不接受这个答案,因为如果您将异常用于预期行为,那么 1)它们不是异常 2)您无法判断异常是否针对您的预期情况而引发或者如果它被抛出是因为您实际上在异步方法中使用了已处置的对象或空引用,而不是套接字。

在带有异步套接字代码的 C++ 中,我会使用 Interlocked 跟踪未完成的异步操作的数量,当需要断开连接时,我会调用 shutdown,然后等待 interlocked 达到 0,然后关闭并销毁任何我需要的成员。

如何在我的 Disconnect 方法中等待所有未完成的异步操作在 C# 中完成以下清单?

using System;
using System.Collections.Generic;
using System.Linq;
using System.Text;
using System.Threading.Tasks;

using log4net;
using System.Net.Sockets;
using System.Net;

namespace IntegrationTests
{
    public class Client2
    {
        class ReceiveContext
        {
            public Socket    m_socket;
            public const int m_bufferSize = 1024;
            public byte[]    m_buffer = new byte[m_bufferSize];
        }

        private static readonly ILog log = LogManager.GetLogger("root");

        static private ulong m_lastId = 1;

        private ulong  m_id;
        private string m_host;
        private uint   m_port;
        private uint   m_timeoutMilliseconds;
        private string m_clientId;
        private Socket m_socket;
        private uint   m_numOutstandingAsyncOps;

        public Client2(string host, uint port, string clientId, uint timeoutMilliseconds)
        {
            m_id                     = m_lastId++;
            m_host                   = host;
            m_port                   = port;
            m_clientId              = clientId;
            m_timeoutMilliseconds    = timeoutMilliseconds;
            m_socket                 = null;
            m_numOutstandingAsyncOps = 0;
        }

        ~Client2()
        {
            Disconnect();
        }

        public void Connect()
        {
            IPHostEntry ipHostInfo = Dns.GetHostEntry(m_host);
            IPAddress[] ipV4Addresses = ipHostInfo.AddressList.Where(x => x.AddressFamily == AddressFamily.InterNetwork).ToArray();
            IPAddress[] ipV6Addresses = ipHostInfo.AddressList.Where(x => x.AddressFamily == AddressFamily.InterNetworkV6).ToArray();
            IPEndPoint endpoint = new IPEndPoint(ipV4Addresses[0], (int)m_port);

            m_socket = new Socket(AddressFamily.InterNetwork, SocketType.Stream, ProtocolType.Tcp);
            m_socket.ReceiveTimeout = (int)m_timeoutMilliseconds;
            m_socket.SendTimeout    = (int)m_timeoutMilliseconds;

            try
            {
                m_socket.Connect(endpoint);

                log.Info(string.Format("Connected to: {0}", m_socket.RemoteEndPoint.ToString()));

                // Issue the next async receive
                ReceiveContext context = new ReceiveContext();
                context.m_socket = m_socket;
                m_socket.BeginReceive(context.m_buffer, 0, ReceiveContext.m_bufferSize, SocketFlags.None, new AsyncCallback(OnReceive), context);
            }
            catch (Exception e)
            {
                // Error
                log.Error(string.Format("Client #{0} Exception caught OnConnect. Exception: {1}"
                                       , m_id, e.ToString()));
            }
        }

        public void Disconnect()
        {
            if (m_socket != null)
            {
                m_socket.Shutdown(SocketShutdown.Both);

                // TODO - <--- Error here in the callbacks where they try to use the socket and it is disposed
                //        We need to wait here until all outstanding async operations complete
                //        Should we use Interlocked to keep track of them and wait on it somehow?
                m_socket.Close();
                m_socket = null;
            }
        }

        public void Login()
        {
            string loginRequest = string.Format("loginstuff{0})", m_clientId);
            var data = Encoding.ASCII.GetBytes(loginRequest);

            m_socket.BeginSend(data, 0, data.Length, 0, new AsyncCallback(OnSend), m_socket);
        }

        public void MakeRequest(string thingy)
        {
            string message = string.Format("requeststuff{0}", thingy);
            var data = Encoding.ASCII.GetBytes(message);

            m_socket.BeginSend(data, 0, data.Length, 0, new AsyncCallback(OnSend), m_socket);
        }

        void OnReceive(IAsyncResult asyncResult)
        {
            ReceiveContext context = (ReceiveContext)asyncResult.AsyncState;

            string data = null;
            try
            {
                int bytesReceived = context.m_socket.EndReceive(asyncResult);
                data = Encoding.ASCII.GetString(context.m_buffer, 0, bytesReceived);

                ReceiveContext newContext = new ReceiveContext();
                newContext.m_socket = context.m_socket;

                m_socket.BeginReceive(newContext.m_buffer, 0, ReceiveContext.m_bufferSize, SocketFlags.None, new AsyncCallback(OnReceive), newContext);
            }
            catch(SocketException e)
            {
                if(e.SocketErrorCode == SocketError.ConnectionAborted) // Check if we disconnected on our end
                {
                    return;
                }
            }
            catch (Exception e)
            {
                // Error
                log.Error(string.Format("Client #{0} Exception caught OnReceive. Exception: {1}"
                                       , m_id, e.ToString()));
            }
        }

        void OnSend(IAsyncResult asyncResult)
        {
            Socket socket = (Socket)asyncResult.AsyncState;

            try
            {
                int bytesSent = socket.EndSend(asyncResult);
            }
            catch(Exception e)
            {
                log.Error(string.Format("Client #{0} Exception caught OnSend. Exception: {1}"
                                       , m_id, e.ToString()));
            }
        }
    }
}

主要:

using System;
using System.Collections.Generic;
using System.Linq;
using System.Text;
using System.Threading.Tasks;

using log4net;
using log4net.Config;

namespace IntegrationTests
{
    class Program
    {
        private static readonly ILog log = LogManager.GetLogger("root");

        static void Main(string[] args)
        {
            try
            {
                XmlConfigurator.Configure();
                log.Info("Starting Component Integration Tests...");

                Client2 client = new Client2("127.0.0.1", 24001, "MyClientId", 60000);
                client.Connect();
                client.Login();
                client.MakeRequest("StuffAndPuff");

                System.Threading.Thread.Sleep(60000); // Sim work until user shutsdown

                client.Disconnect();
            }
            catch (Exception e)
            {
                log.Error(string.Format("Caught an exception in main. Exception: {0}"
                                      , e.ToString()));
            }
        }
    }
}

编辑:

这是我尽我所能使用 Evk 提出的答案的额外尝试。据我所知,它工作正常。

问题在于,我觉得我基本上将所有异步调用都变成了同步调用,因为需要锁定任何会改变计数器或套接字状态的东西。同样,与我的 C++ 相比,我是 C# 的新手,所以请指出我是否完全错过了解释他的答案的标记。

using System;
using System.Collections.Generic;
using System.Linq;
using System.Net;
using System.Net.Sockets;
using System.Text;
using System.Threading;
using System.Threading.Tasks;

namespace IntegrationTests
{
    public class Client
    {
        class ReceiveContext
        {
            public const int     _bufferSize    = 1024;
            public byte[]        _buffer        = new byte[_bufferSize]; // Contains bytes from one receive
            public StringBuilder _stringBuilder = new StringBuilder();   // Contains bytes for multiple receives in order to build message up to delim
        }

        private static readonly ILog _log = LogManager.GetLogger("root");

        static private ulong _lastId = 1;
        private ulong  _id;

        protected string         _host;
        protected int            _port;
        protected int            _timeoutMilliseconds;
        protected string         _sessionId;
        protected Socket         _socket;
        protected object         _lockNumOutstandingAsyncOps;
        protected int            _numOutstandingAsyncOps;
        private bool             _disposed = false;

        public Client(string host, int port, string sessionId, int timeoutMilliseconds)
        {
            _id                         = _lastId++;
            _host                       = host;
            _port                       = port;
            _sessionId                  = sessionId;
            _timeoutMilliseconds        = timeoutMilliseconds;
            _socket                     = null;
            _numOutstandingAsyncOps     = 0;
            _lockNumOutstandingAsyncOps = new object();
        }

        public void Dispose()
        {
            Dispose(true);
            GC.SuppressFinalize(this);
        }

        protected virtual void Dispose(bool disposing)
        {
            if(_disposed)
            {
                return;
            }

            if (disposing)
            {
                _socket.Close();
            }

            _disposed = true;
        }

        public void Connect()
        {
            lock (_lockNumOutstandingAsyncOps)
            {
                IPHostEntry ipHostInfo = Dns.GetHostEntry(_host);
                IPAddress[] ipV4Addresses = ipHostInfo.AddressList.Where(x => x.AddressFamily == AddressFamily.InterNetwork).ToArray();
                IPAddress[] ipV6Addresses = ipHostInfo.AddressList.Where(x => x.AddressFamily == AddressFamily.InterNetworkV6).ToArray();
                IPEndPoint endpoint = new IPEndPoint(ipV4Addresses[0], _port);

                _socket = new Socket(AddressFamily.InterNetwork, SocketType.Stream, ProtocolType.Tcp);
                _socket.ReceiveTimeout = _timeoutMilliseconds;
                _socket.SendTimeout = _timeoutMilliseconds;

                try
                {
                    _socket.Connect(endpoint);
                }
                catch (Exception e)
                {
                    // Error
                    Debug.WriteLine(string.Format("Client #{0} Exception caught OnConnect. Exception: {1}"
                                           , _id, e.ToString()));
                    return;
                }

                Debug.WriteLine(string.Format("Client #{0} connected to: {1}", _id, _socket.RemoteEndPoint.ToString()));

                // Issue the first async receive
                ReceiveContext context = new ReceiveContext();

                ++_numOutstandingAsyncOps;
                _socket.BeginReceive(context._buffer, 0, ReceiveContext._bufferSize, SocketFlags.None, new AsyncCallback(OnReceive), context);
            }
        }

        public void Disconnect()
        {
            if (_socket != null)
            {
                // We need to wait here until all outstanding async operations complete
                // In order to avoid getting 'Object was disposed' exceptions in those async ops that use the socket
                lock(_lockNumOutstandingAsyncOps)
                {
                    Debug.WriteLine(string.Format("Client #{0} Disconnecting...", _id));

                    _socket.Shutdown(SocketShutdown.Both);

                    while (_numOutstandingAsyncOps > 0)
                    {
                        Monitor.Wait(_lockNumOutstandingAsyncOps);
                    }

                    _socket.Close();
                    _socket = null;
                }
            }
        }

        public void Login()
        {
            lock (_lockNumOutstandingAsyncOps)
            {
                if (_socket != null && _socket.Connected)
                {
                    string loginRequest = string.Format("loginstuff{0}", _clientId);
                    var data = Encoding.ASCII.GetBytes(loginRequest);

                    Debug.WriteLine(string.Format("Client #{0} Sending Login Request: {1}"
                                           , _id, loginRequest));

                    ++_numOutstandingAsyncOps;
                    _socket.BeginSend(data, 0, data.Length, 0, new AsyncCallback(OnSend), _socket);
                }
                else
                {
                    Debug.WriteLine(string.Format("Client #{0} Login was called, but Socket is null or no longer connected."
                                           , _id));
                }
            }
        }

        public void MakeRequest(string thingy)
        {
            lock (_lockNumOutstandingAsyncOps)
            {
                if (_socket != null && _socket.Connected)
                {
                    string message = string.Format("requeststuff{0}", thingy);
                    var data = Encoding.ASCII.GetBytes(message);

                    Debug.WriteLine(string.Format("Client #{0} Sending Request: {1}"
                                           , _id, message));

                    ++_numOutstandingAsyncOps;
                    _socket.BeginSend(data, 0, data.Length, 0, new AsyncCallback(OnSend), _socket);
                }
                else
                {
                    Debug.WriteLine(string.Format("Client #{0} MakeRequest was called, but Socket is null or no longer connected."
                                           , _id));
                }
            }
        }

        protected void OnReceive(IAsyncResult asyncResult)
        {
            lock (_lockNumOutstandingAsyncOps)
            {
                ReceiveContext context = (ReceiveContext)asyncResult.AsyncState;

                string data = null;

                try
                {
                    int bytesReceived = _socket.EndReceive(asyncResult);
                    data = Encoding.ASCII.GetString(context._buffer, 0, bytesReceived);

                    // If the remote host shuts down the Socket connection with the Shutdown method, and all available data has been received,
                    // the EndReceive method will complete immediately and return zero bytes
                    if (bytesReceived > 0)
                    {
                        StringBuilder stringBuilder = context._stringBuilder.Append(data);

                        int index = -1;
                        do
                        {
                            index = stringBuilder.ToString().IndexOf("#");
                            if (index != -1)
                            {
                                string message = stringBuilder.ToString().Substring(0, index + 1);
                                stringBuilder.Remove(0, index + 1);

                                Debug.WriteLine(string.Format("Client #{0} Received Data: {1}"
                                                       , _id, message));
                            }
                        } while (index != -1);
                    }
                }
                catch (SocketException e)
                {
                    // Check if we disconnected on our end
                    if (e.SocketErrorCode == SocketError.ConnectionAborted)
                    {
                        // Ignore
                    }
                    else
                    {
                        // Error
                        Debug.WriteLine(string.Format("Client #{0} SocketException caught OnReceive. Exception: {1}"
                                               , _id, e.ToString()));
                        Disconnect();
                    }
                }
                catch (Exception e)
                {
                    // Error
                    Debug.WriteLine(string.Format("Client #{0} Exception caught OnReceive. Exception: {1}"
                                           , _id, e.ToString()));
                    Disconnect();
                }
                finally
                {
                    --_numOutstandingAsyncOps;
                    Monitor.Pulse(_lockNumOutstandingAsyncOps);
                }
            }

            // Issue the next async receive
            lock (_lockNumOutstandingAsyncOps)
            {
                if (_socket != null && _socket.Connected)
                {
                    ++_numOutstandingAsyncOps;

                    ReceiveContext newContext = new ReceiveContext();
                    _socket.BeginReceive(newContext._buffer, 0, ReceiveContext._bufferSize, SocketFlags.None, new AsyncCallback(OnReceive), newContext);
                }
            }
        }

        protected void OnSend(IAsyncResult asyncResult)
        {
            lock (_lockNumOutstandingAsyncOps)
            {
                try
                {
                    int bytesSent = _socket.EndSend(asyncResult);
                }
                catch (Exception e)
                {
                    Debug.WriteLine(string.Format("Client #{0} Exception caught OnSend. Exception: {1}"
                                           , _id, e.ToString()));
                    Disconnect();
                }
                finally
                {
                    --_numOutstandingAsyncOps;
                    Monitor.Pulse(_lockNumOutstandingAsyncOps);
                }
            }
        }
    }
}

【问题讨论】:

  • 请注意,终结器用于清理非托管资源,而不是托管资源。你的对象根本不应该有终结器。
  • CountdownEvent 是你的朋友。
  • 补充一下 Servy 所说的,因为您来自 C++ 背景,所以不要将 C++ 析构函数与托管终结器混淆。事实上,C# 借用了语法,甚至在某些文档中使用了“析构函数”这个词,这真的很不幸。 IDisposable 是 C# 用于确定性清理的东西,并且很少需要终结器(当对象被垃圾回收时,它们是释放非托管资源的最后权宜之计——通常,这样的对象应该已经调用了它的 Dispose 并且终结器被抑制)。
  • 关于一个反对意见:“您无法判断异常是针对您的预期情况引发的,还是因为您实际上在异步方法中使用了已处置的对象或空引用而不是套接字而引发的。”目的是您only 捕获ObjectDisposedException,并且only 捕获Socket.EndReceive。当且仅当套接字已关闭时,才会抛出 ObjectDisposedException。块中不应有其他语句。如果你将一堆东西包装在一个通用的 catch (Exception) 中并设置对 null 的引用,那么是的,你有问题。
  • @JeroenMostert 也许,但是 CountdownEvent 不能从 0 开始,因为它会被认为是有信号的,这很好,除了你不能在那之后添加计数。所以,我不知道我会如何使用它。这是我用于测试您的建议的概念程序:ideone.com/zfbJH9

标签: c# sockets asynchronous synchronization interlocked


【解决方案1】:

您可以为此使用Monitor.WaitMonitor.Pulse

static int _outstandingOperations;
static readonly object _lock = new object();
static void Main() {
    for (int i = 0; i < 100; i++) {
        var tmp = i;
        Task.Run(() =>
        {
            lock (_lock) {
                _outstandingOperations++;
            }
            // some work
            Thread.Sleep(new Random(tmp).Next(0, 5000));
            lock (_lock) {
                _outstandingOperations--;
                // notify condition might have changed
                Monitor.Pulse(_lock);
            }
        });
    }

    lock (_lock) {
        // condition check
        while (_outstandingOperations > 0)
            // will wait here until pulsed, lock will be released during wait
            Monitor.Wait(_lock);
    }
}

【讨论】:

  • 打算编辑原始帖子并尝试与此相关的问题。
  • @ChristopherPisz 您不需要像在示例中那样锁定整个操作,只需要围绕递增/递减字段(就像在我的回答中一样 - 我不锁定线程睡眠模仿工作) .您可以将其移至单独的方法(IncAnynsOperations/DecAsyncOperations)。
  • 我只锁定了计数的递增和递减,并且在调用断开连接时得到了相同的旧异常,因为异步方法仍在作用于套接字的中间,所以套接字状态必须是也被锁定了。
  • @ChristopherPisz 在您提供给您的代码中 first 关闭套接字,然后才等待异步操作降至 0。当然这不起作用。您必须首先等待它们降至 0,然后关闭(最好使用一些标志指示正在关闭,这样您就不会在它们降至 0 后启动更多异步操作)。
  • 是的,我认为它需要一些状态跟踪。今天第三次尝试。
猜你喜欢
  • 2012-04-27
  • 2017-02-15
  • 1970-01-01
  • 2011-01-24
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多