Skip to content

Commit 18e4ede

Browse files
committed
Improve handling of connections closed by the peer (#248)
***NO_CI*** (cherry picked from commit 6803587)
1 parent ca8a034 commit 18e4ede

3 files changed

Lines changed: 167 additions & 100 deletions

File tree

‎WebSockets/ReceiveAndControllThread.cs‎

Lines changed: 124 additions & 89 deletions
Original file line numberDiff line numberDiff line change
@@ -26,98 +26,124 @@ public ReceiveAndControllThread(WebSocket webSocket)
2626

2727
while (!_webSocket.Stopped)
2828
{
29-
var messageFrame = _webSocket.WebSocketReceiver.StartReceivingMessage();
30-
31-
if (messageFrame == null)
29+
try
3230
{
33-
//Here we could let the thread sleep to safe resources
31+
ProcessIncomingMessage();
3432
}
35-
else
33+
catch (Exception ex)
3634
{
37-
//handle error
38-
if (messageFrame.Error)
35+
// failure reading from the stream (connection closed, reset, disposed, etc.)
36+
if (!_webSocket.Stopped)
3937
{
4038
_webSocket.HasError = true;
4139

42-
Debug.WriteLine($"{_webSocket.RemoteEndPoint} closed with error: {messageFrame.ErrorMessage}");
40+
Debug.WriteLine($"{_webSocket.RemoteEndPoint} closed with error: {ex.Message}");
4341

44-
_webSocket.RawClose(messageFrame.CloseStatus, Encoding.UTF8.GetBytes(messageFrame.ErrorMessage), true);
42+
// don't leak internal exception details to the peer
43+
_webSocket.RawClose(WebSocketCloseStatus.EndpointUnavailable, Encoding.UTF8.GetBytes("Connection error"), true);
44+
45+
// RawClose returns without closing if a close was already sent (or the socket is closing),
46+
// and the timeout checker is disposed below, so make sure the connection is closed
47+
_webSocket.HardClose();
4548
}
46-
else if (messageFrame.IsControllFrame)
47-
{
48-
byte[] buffer = _webSocket.WebSocketReceiver.ReadBuffer(messageFrame.MessageLength, messageFrame.Masks);
4949

50-
_webSocket.LastContactTimeStamp = DateTime.UtcNow;
50+
// stream can't be used anymore
51+
break;
52+
}
53+
}
5154

52-
switch (messageFrame.OpCode)
53-
{
54-
case OpCode.PingFrame:
55-
// need to send a Pong
56-
var pong = new SendMessageFrame() { Buffer = buffer, OpCode = OpCode.PongFrame};
57-
messageFrame.OpCode = OpCode.PongFrame;
58-
messageFrame.IsMasked = false;
59-
_webSocket.QueueMessageToSend(pong);
60-
break;
61-
62-
case OpCode.PongFrame:
63-
// received a Pong
64-
// checking if content Pong matches Ping is not implemented due to thread safety and memory consumption considerations
65-
_webSocket.Pinging = false;
66-
break;
67-
68-
case OpCode.ConnectionCloseFrame:
69-
_webSocket.CloseStatus = WebSocketCloseStatus.Empty;
70-
71-
if (buffer.Length > 1)
72-
{
73-
byte[] closeByteCode = new byte[] { buffer[1], buffer[0] };
74-
UInt16 statusCode = BitConverter.ToUInt16(closeByteCode, 0);
75-
if (statusCode > 999 && statusCode < 1012)
76-
{
77-
_webSocket.CloseStatus = (WebSocketCloseStatus)statusCode;
78-
}
79-
}
55+
_webSocket.ReceiveStream.Close();
56+
timeoutCheckerTimer.Change(Timeout.Infinite, Timeout.Infinite);
57+
timeoutCheckerTimer.Dispose();
58+
}
8059

81-
//connection asked to be closed return answer
82-
if (_webSocket.State != WebSocketFrame.WebSocketState.CloseSent)
83-
{
84-
_webSocket.State = WebSocketFrame.WebSocketState.CloseReceived;
60+
private void ProcessIncomingMessage()
61+
{
62+
var messageFrame = _webSocket.WebSocketReceiver.StartReceivingMessage();
8563

86-
_webSocket.RawClose(WebSocketCloseStatus.NormalClosure, buffer, true);
87-
}
88-
//response to connection close we can shut down the socket.
89-
else
64+
if (messageFrame == null)
65+
{
66+
//Here we could let the thread sleep to safe resources
67+
}
68+
else
69+
{
70+
//handle error
71+
if (messageFrame.Error)
72+
{
73+
_webSocket.HasError = true;
74+
75+
Debug.WriteLine($"{_webSocket.RemoteEndPoint} closed with error: {messageFrame.ErrorMessage}");
76+
77+
_webSocket.RawClose(messageFrame.CloseStatus, Encoding.UTF8.GetBytes(messageFrame.ErrorMessage), true);
78+
}
79+
else if (messageFrame.IsControllFrame)
80+
{
81+
byte[] buffer = _webSocket.WebSocketReceiver.ReadBuffer(messageFrame.MessageLength, messageFrame.Masks);
82+
83+
_webSocket.LastContactTimeStamp = DateTime.UtcNow;
84+
85+
switch (messageFrame.OpCode)
86+
{
87+
case OpCode.PingFrame:
88+
// need to send a Pong
89+
var pong = new SendMessageFrame() { Buffer = buffer, OpCode = OpCode.PongFrame};
90+
messageFrame.OpCode = OpCode.PongFrame;
91+
messageFrame.IsMasked = false;
92+
_webSocket.QueueMessageToSend(pong);
93+
break;
94+
95+
case OpCode.PongFrame:
96+
// received a Pong
97+
// checking if content Pong matches Ping is not implemented due to thread safety and memory consumption considerations
98+
_webSocket.Pinging = false;
99+
break;
100+
101+
case OpCode.ConnectionCloseFrame:
102+
_webSocket.CloseStatus = WebSocketCloseStatus.Empty;
103+
104+
if (buffer.Length > 1)
105+
{
106+
byte[] closeByteCode = new byte[] { buffer[1], buffer[0] };
107+
UInt16 statusCode = BitConverter.ToUInt16(closeByteCode, 0);
108+
if (statusCode > 999 && statusCode < 1012)
90109
{
91-
_webSocket.HardClose();
110+
_webSocket.CloseStatus = (WebSocketCloseStatus)statusCode;
92111
}
93-
break;
94-
}
112+
}
113+
114+
//connection asked to be closed return answer
115+
if (_webSocket.State != WebSocketFrame.WebSocketState.CloseSent)
116+
{
117+
_webSocket.State = WebSocketFrame.WebSocketState.CloseReceived;
118+
119+
_webSocket.RawClose(WebSocketCloseStatus.NormalClosure, buffer, true);
120+
}
121+
//response to connection close we can shut down the socket.
122+
else
123+
{
124+
_webSocket.HardClose();
125+
}
126+
break;
95127
}
96-
else
128+
}
129+
else
130+
{
131+
if (messageFrame.Error)
97132
{
98-
if (messageFrame.Error)
99-
{
100-
Debug.WriteLine($"Error message from '{_webSocket.RemoteEndPoint}' error - {messageFrame.ErrorMessage}");
133+
Debug.WriteLine($"Error message from '{_webSocket.RemoteEndPoint}' error - {messageFrame.ErrorMessage}");
101134

102-
_webSocket.RawClose(messageFrame.CloseStatus, Encoding.UTF8.GetBytes(messageFrame.ErrorMessage), true);
103-
}
104-
else
105-
{
106-
messageFrame.Buffer = _webSocket.WebSocketReceiver.ReadBuffer(messageFrame.MessageLength, messageFrame.Masks);
135+
_webSocket.RawClose(messageFrame.CloseStatus, Encoding.UTF8.GetBytes(messageFrame.ErrorMessage), true);
136+
}
137+
else
138+
{
139+
messageFrame.Buffer = _webSocket.WebSocketReceiver.ReadBuffer(messageFrame.MessageLength, messageFrame.Masks);
107140

108-
_webSocket.LastContactTimeStamp = DateTime.UtcNow;
141+
_webSocket.LastContactTimeStamp = DateTime.UtcNow;
109142

110-
OnNewMessage(messageFrame);
111-
}
143+
OnNewMessage(messageFrame);
112144
}
113145
}
114-
115-
116146
}
117-
118-
_webSocket.ReceiveStream.Close();
119-
timeoutCheckerTimer.Change(Timeout.Infinite, Timeout.Infinite);
120-
timeoutCheckerTimer.Dispose();
121147
}
122148

123149

@@ -129,33 +155,42 @@ private void CheckTimeouts(object thread)
129155
receiveThread.Suspend();
130156
#pragma warning restore S3889 // Neither "Thread.Resume" nor "Thread.Suspend" should be used
131157

132-
//Controlling ping and ControllerMessagesTimeout
133-
if (_webSocket.Pinging
134-
&& _webSocket.PingTime.Add(_webSocket.ServerTimeout) < DateTime.UtcNow)
158+
try
135159
{
136-
_webSocket.RawClose(WebSocketCloseStatus.PolicyViolation, Encoding.UTF8.GetBytes("Ping timeout"), true);
160+
//Controlling ping and ControllerMessagesTimeout
161+
if (_webSocket.Pinging
162+
&& _webSocket.PingTime.Add(_webSocket.ServerTimeout) < DateTime.UtcNow)
163+
{
164+
_webSocket.RawClose(WebSocketCloseStatus.PolicyViolation, Encoding.UTF8.GetBytes("Ping timeout"), true);
137165

138-
Debug.WriteLine($"{_webSocket.RemoteEndPoint} ping timed out");
139-
}
166+
Debug.WriteLine($"{_webSocket.RemoteEndPoint} ping timed out");
167+
}
140168

141-
if (_webSocket.State == WebSocketFrame.WebSocketState.CloseSent
142-
&& _webSocket.ClosingTime.Add(_webSocket.ServerTimeout) < DateTime.UtcNow)
143-
{
144-
_webSocket.HardClose();
145-
}
169+
if (_webSocket.State == WebSocketFrame.WebSocketState.CloseSent
170+
&& _webSocket.ClosingTime.Add(_webSocket.ServerTimeout) < DateTime.UtcNow)
171+
{
172+
_webSocket.HardClose();
173+
}
146174

147-
if (_webSocket.KeepAliveInterval != Timeout.InfiniteTimeSpan
148-
&& _webSocket.State != WebSocketFrame.WebSocketState.CloseSent
149-
&& !_webSocket.Pinging
150-
&& _webSocket.LastContactTimeStamp.Add(_webSocket.KeepAliveInterval) < DateTime.UtcNow)
175+
if (_webSocket.KeepAliveInterval != Timeout.InfiniteTimeSpan
176+
&& _webSocket.State != WebSocketFrame.WebSocketState.CloseSent
177+
&& !_webSocket.Pinging
178+
&& _webSocket.LastContactTimeStamp.Add(_webSocket.KeepAliveInterval) < DateTime.UtcNow)
179+
{
180+
_webSocket.SendPing();
181+
}
182+
}
183+
catch (Exception ex)
151184
{
152-
_webSocket.SendPing();
185+
Debug.WriteLine($"{_webSocket.RemoteEndPoint} error checking timeouts: {ex.Message}");
153186
}
154-
187+
finally
188+
{
189+
// always resume the receive thread, otherwise it will be left suspended forever
155190
#pragma warning disable S3889 // OK to use in .NET nanoFramework context
156-
receiveThread.Resume();
191+
receiveThread.Resume();
157192
#pragma warning restore S3889 // Neither "Thread.Resume" nor "Thread.Suspend" should be used
158-
193+
}
159194
}
160195

161196
private void OnNewMessage(ReceiveMessageFrame message)

‎WebSockets/WebSocket.cs‎

Lines changed: 36 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -27,6 +27,7 @@ public abstract class WebSocket : IDisposable
2727

2828
internal WebSocketSender _webSocketSender;
2929
private readonly object _syncLock = new object();
30+
private int _hardClosed = 0;
3031
internal MessageReceivedEventHandler CallbacksMessageReceivedEventHandler;
3132

3233
/// <summary>
@@ -133,6 +134,7 @@ protected void ConnectToStream(NetworkStream stream, bool isServer, Socket socke
133134
_socket = socket;
134135
RemoteEndPoint = (IPEndPoint)socket.RemoteEndPoint;
135136
LastContactTimeStamp = DateTime.UtcNow;
137+
_hardClosed = 0;
136138

137139
//start server sending and receiving async
138140
WebSocketReceiver = new WebSocketReceiver(stream, RemoteEndPoint, this, IsServer, MaxReceiveFrameSize, OnReadError);
@@ -275,15 +277,26 @@ internal void RawClose(WebSocketCloseStatus closeStatus = WebSocketCloseStatus.E
275277

276278
if (CloseImmediately)
277279
{
280+
State = WebSocketState.CloseSent;
281+
ClosingTime = DateTime.UtcNow;
282+
283+
// Give it a moment for sending a close message. This will block the thread.
284+
// Wait is bounded because the send thread can be stuck on a dead connection.
285+
int maxWaitMs = (int)ServerTimeout.TotalMilliseconds;
286+
287+
if (maxWaitMs <= 0)
288+
{
289+
// infinite (or invalid) server timeout, use a sensible default
290+
maxWaitMs = 5000;
291+
}
292+
278293
int msWaited = 0;
279294

280-
//Give it a moment for sending a close message. This will block the thread.
281-
while (!_webSocketSender.CloseMessageSent )
295+
while (!_webSocketSender.CloseMessageSent
296+
&& msWaited < maxWaitMs)
282297
{
283298
msWaited += 50;
284299
Thread.Sleep(50);
285-
State = WebSocketState.CloseSent;
286-
ClosingTime = DateTime.UtcNow;
287300
}
288301

289302
HardClose();
@@ -292,24 +305,36 @@ internal void RawClose(WebSocketCloseStatus closeStatus = WebSocketCloseStatus.E
292305

293306
internal void HardClose()
294307
{
308+
// can be called from several threads (receive, send error, timeout checker)
309+
// lock-free guard: the timeout checker suspends the receive thread, so a lock here could deadlock
310+
if (Interlocked.CompareExchange(ref _hardClosed, 1, 0) != 0)
311+
{
312+
return;
313+
}
314+
295315
State = WebSocketState.Closed;
296316

297317
StopReceiving();
298318
_webSocketSender.StopSender();
299319

300320
Debug.WriteLine($"Connection - {RemoteEndPoint.ToString()} - Closed");
301-
302-
ConnectionClosed?.Invoke(this, new EventArgs());
303321

304-
//Let the tcp socket linger for a second so it can try and send all data out before final close.
305322
try
306323
{
307-
_socket.SetSocketOption(SocketOptionLevel.Socket, SocketOptionName.Linger, 1);
308-
_socket.Close();
324+
ConnectionClosed?.Invoke(this, new EventArgs());
309325
}
310-
catch (ObjectDisposedException e)
326+
finally
311327
{
312-
Debug.WriteLine("socket could not be closed properly because it was already disposed");
328+
//Let the tcp socket linger for a second so it can try and send all data out before final close.
329+
try
330+
{
331+
_socket.SetSocketOption(SocketOptionLevel.Socket, SocketOptionName.Linger, 1);
332+
_socket.Close();
333+
}
334+
catch (ObjectDisposedException e)
335+
{
336+
Debug.WriteLine("socket could not be closed properly because it was already disposed");
337+
}
313338
}
314339
}
315340

‎WebSockets/WebSocketReceiver.cs‎

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -199,6 +199,13 @@ byte[] ReadFixedSizeBuffer(int size, byte[] masks = null)
199199
while (size > 0)
200200
{
201201
int bytes = _inputStream.Read(buffer, offset, size);
202+
203+
if (bytes <= 0)
204+
{
205+
// zero bytes read means the remote end has closed the connection
206+
throw new SocketException(SocketError.ConnectionReset);
207+
}
208+
202209
offset += bytes;
203210
size -= bytes;
204211
}

0 commit comments

Comments
 (0)