-
Notifications
You must be signed in to change notification settings - Fork 3
Expand file tree
/
Copy pathClient.Heartbeat.cs
More file actions
139 lines (115 loc) · 4.61 KB
/
Copy pathClient.Heartbeat.cs
File metadata and controls
139 lines (115 loc) · 4.61 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
using System;
using System.Collections.Concurrent;
using System.Linq;
using System.Threading.Tasks;
using GoPlay.Core;
using GoPlay.Core.Protocols;
using GoPlay.Core.Utils;
using GoPlay.Exceptions;
namespace GoPlay
{
public partial class Client<T>
{
protected Task m_heartbeatTask;
protected ConcurrentDictionary<uint, DateTime> m_pingDict = new ConcurrentDictionary<uint, DateTime>();
protected int m_duration;
protected int m_pingCount;
protected TimeSpan m_pingAvg = TimeSpan.Zero;
protected TimeSpan m_pingMax = TimeSpan.MinValue;
protected TimeSpan m_pingMin = TimeSpan.MaxValue;
public Task HeartbeatTask => m_heartbeatTask;
public int PingCount => m_pingCount;
public TimeSpan PingAvg => m_pingAvg;
public TimeSpan PingMax => m_pingMax;
public TimeSpan PingMin => m_pingMin;
protected virtual void StartHeartbeat()
{
m_heartbeatTask = TaskUtil.LongRun(HeartbeatLoop, m_cancelSource.Token);
}
protected virtual void HeartbeatLoop()
{
while (!m_cancelSource.Token.IsCancellationRequested)
{
try
{
var task = Task.Delay(Consts.HeartBeat.Update);
task.Wait(m_cancelSource.Token);
if (task.IsCanceled) return;
if (!IsConnected) continue;
m_duration -= (int)Consts.HeartBeat.Update.TotalMilliseconds;
//check time out
var isTimeOut = false;
var dict = m_pingDict.ToList();
foreach (var kv in dict)
{
var ts = DateTime.UtcNow.Subtract(kv.Value);
if (ts < Consts.HeartBeat.Timeout) continue;
isTimeOut = true;
ResetHeartBeatData();
OnErrorEvent(new HeartbeatTimeoutException());
DisconnectAsync().ConfigureAwait(false);
break;
}
if (isTimeOut) return;
if (m_duration > 0) continue;
m_duration = (int)m_handshake.HeartBeatInterval;
if (!m_pingDict.IsEmpty) continue;
//calculate ping timeout
var pack = Package.Create(0, PackageType.Ping, EncodingType);
m_pingDict.TryAdd(pack.Header.PackageInfo.Id, DateTime.UtcNow);
Send(pack);
}
catch (OperationCanceledException)
{
//IGNORE ERR
}
catch (AggregateException err)
{
if (err.InnerException is OperationCanceledException) continue;
if (err.InnerException is TaskCanceledException) continue;
OnErrorEvent(err);
DisconnectAsync().ConfigureAwait(false);
}
catch (Exception err)
{
OnErrorEvent(err);
DisconnectAsync().ConfigureAwait(false);
}
}
}
protected virtual void ResetHeartBeatData()
{
m_pingDict.Clear();
m_duration = 0;
m_pingCount = 0;
m_pingAvg = TimeSpan.Zero;
m_pingMax = TimeSpan.MinValue;
m_pingMin = TimeSpan.MaxValue;
}
protected virtual void ResolvePing(Package pack)
{
var resp = new Package
{
Header = pack.Header.Clone(),
};
resp.Header.PackageInfo.Type = PackageType.Pong;
Send(resp);
}
protected virtual void ResolvePong(Package pack)
{
if (!m_pingDict.TryRemove(pack.Header.PackageInfo.Id, out var dateTime)) return;
m_pingCount++;
var ts = DateTime.UtcNow.Subtract(dateTime);
if (ts < m_pingMin) m_pingMin = ts;
if (ts > m_pingMax) m_pingMax = ts;
var avgTick = (double)(m_pingCount - 1) / m_pingCount * m_pingAvg.Ticks;
avgTick += (double)ts.Ticks / m_pingCount;
m_pingAvg = TimeSpan.FromTicks((long)avgTick);
}
//TODO: remove test function
public string NetworkStatus()
{
return $"Ping: avg={m_pingAvg.TotalMilliseconds} ms, max={m_pingMax.TotalMilliseconds} ms, min={m_pingMin.TotalMilliseconds} ms";
}
}
}