-
Notifications
You must be signed in to change notification settings - Fork 391
Expand file tree
/
Copy pathLogicService.LogicThread.cs
More file actions
442 lines (380 loc) · 22.3 KB
/
Copy pathLogicService.LogicThread.cs
File metadata and controls
442 lines (380 loc) · 22.3 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
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
using Intersect.Server.Database;
using Intersect.Server.General;
using Intersect.Server.Maps;
using Intersect.Threading;
using System.Diagnostics;
using Intersect.Server.Entities;
using Amib.Threading;
using System.Collections.Concurrent;
using Intersect.Core;
using Intersect.Framework.Core;
using Intersect.Server.Metrics;
using Intersect.Server.Networking;
using Intersect.Server.Database.PlayerData.Players;
using Intersect.Server.Database.PlayerData.Api;
using Intersect.Server.Core.MapInstancing;
namespace Intersect.Server.Core;
internal sealed partial class LogicService
{
internal static INetworkPoolDataProvider NetworkPoolDataProvider { get; set; } =
new SinglePlayerNetworkPoolDataProvider();
internal sealed partial class LogicThread : Threaded<ServerContext>
{
private readonly LogicService _logicService;
private long _nextClearExpiredTokens;
/// <summary>
/// We lock on this in order to stop maps from entering the update queue. This is only done when the editor is saving/modifying game maps or the map grids are being rebuilt.
/// </summary>
public readonly object LogicLock = new object();
/// <summary>
/// This is our thread pool for handling server/game logic. This includes npcs, event processing, map updating, projectiles, spell casting, etc.
/// Min/Max Number of Threads & Idle Timeouts are set via server config.
/// </summary>
public readonly SmartThreadPool LogicPool = new SmartThreadPool(
new STPStartInfo()
{
ThreadPoolName = "LogicPool",
IdleTimeout = 20000,
MinWorkerThreads = Options.Instance.Processing.MinLogicThreads,
MaxWorkerThreads = Options.Instance.Processing.MaxLogicThreads
}
);
public LogicThread(LogicService logicService) : base("ServerLogic")
{
_logicService = logicService;
}
/// <summary>
/// Queue of active maps which maps are added to after being updated. Once a map makes it to the front of the queue they are updated again.
/// </summary>
public readonly ConcurrentQueue<MapInstance> MapInstanceUpdateQueue = new ConcurrentQueue<MapInstance>();
/// <summary>
/// Queue of active maps which maps are added to after being updated. Once a map makes it to the front of the queue they are updated again.
/// This queue is only used for projectile updates if the projectile update interval does not match the map update interval in the server config.
/// </summary>
public readonly ConcurrentQueue<MapInstance> MapInstanceProjectileQueue = new ConcurrentQueue<MapInstance>();
/// <summary>
/// This is the set of maps determined to be 'active' based on player locations in the game. Our logic recalculates this hashset every 250ms.
/// When maps are updated they are not added back into the map update queues unless they exist in this hash set.
/// </summary>
public readonly Dictionary<Guid, MapInstance> ActiveMapInstances = new Dictionary<Guid, MapInstance>();
// Reusable buffer for UpdateInstanceControllers — avoids ToArray() every 250ms.
private MapInstance[] _activeMapInstancesBuffer = Array.Empty<MapInstance>();
// Track last console title to avoid redundant syscalls every second.
private string _lastConsoleTitle = string.Empty;
protected override void ThreadStart(ServerContext serverContext)
{
if (serverContext == null)
{
throw new ArgumentNullException(nameof(serverContext));
}
try
{
var swCpsTimer = Timing.Global.Milliseconds + 1000;
var lastCpuTime = Process.GetCurrentProcess().TotalProcessorTime;
var saveServerVariablesTimer = Timing.Global.Milliseconds + Options.Instance.Processing.DatabaseSaveServerVariablesInterval;
var metricsTimer = 0l;
long swCps = 0;
long updateTimer = 0;
var processedMapInstances = new HashSet<Guid>();
var sourceMapInstances = new HashSet<Guid>();
var toRemove = new List<Guid>();
var players = 0;
while (ServerContext.Instance.IsRunning)
{
var startTime = Timing.Global.Milliseconds;
// Cron-clear expired refresh tokens
if (startTime > _nextClearExpiredTokens)
{
#pragma warning disable CA2008 // Do not create tasks without passing a TaskScheduler
_ = RefreshToken
.RemoveExpiredAsync(250)
/* Make it speed up the next call if we definitely have more expired tokens */
.ContinueWith(remainingCount => _nextClearExpiredTokens -= remainingCount.Result > 0 ? 60000 : 0, TaskContinuationOptions.RunContinuationsAsynchronously);
#pragma warning restore CA2008 // Do not create tasks without passing a TaskScheduler
_nextClearExpiredTokens = startTime + 60000;
}
if (startTime > updateTimer)
{
//Resync Active Maps By Scanning Players and Their Surrounding Maps
players = 0;
processedMapInstances.Clear();
sourceMapInstances.Clear();
//Metrics
var globalEntities = 0;
var events = 0;
var eventsProcessing = 0;
var autorunEvents = 0;
var onlinePlayers = Player.OnlinePlayers;
foreach (var player in onlinePlayers)
{
if (player != null)
{
players++;
if (Options.Instance.Metrics.Enable)
{
events += player.EventLookup.Count;
foreach (var evt in player.EventLookup.Values)
{
if (evt.CallStack?.Count > 0) eventsProcessing++;
}
autorunEvents += player.CommonAutorunEvents + player.MapAutorunEvents;
}
}
var plyrMap = player?.MapId ?? Guid.Empty;
if (plyrMap != Guid.Empty
&& !sourceMapInstances.Contains(plyrMap))
{
// Queue up each surrounding map instance of the given player
MapController.GetSurroundingMapInstances(plyrMap, player.MapInstanceId, true)
.ForEach(instance =>
{
if (!processedMapInstances.Contains(instance.Id))
{
if (!ActiveMapInstances.ContainsKey(instance.Id)) // ContainsKey() directly instead of .Keys.Contains().
{
AddToQueue(instance);
}
globalEntities += instance.GetCachedEntities().Length;
processedMapInstances.Add(instance.Id);
}
sourceMapInstances.Add(instance.Id);
});
}
}
// Iterate dictionary directly; collect removals in reusable list.
toRemove.Clear();
foreach (var (instanceId, mapInstance) in ActiveMapInstances)
{
if (processedMapInstances.Contains(instanceId) || mapInstance.ShouldBeActive())
{
continue;
}
if (mapInstance != default)
{
if (mapInstance.ShouldBeCleaned())
{
mapInstance.RemoveLayerFromController();
}
else if (!mapInstance.ShouldBeActive())
{
mapInstance.ResetNpcSpawns();
}
}
toRemove.Add(instanceId);
}
foreach (var id in toRemove)
{
ActiveMapInstances.Remove(id);
}
// Reuse buffer; only reallocate when count changes.
var activeCount = ActiveMapInstances.Count;
if (_activeMapInstancesBuffer.Length != activeCount)
{
_activeMapInstancesBuffer = new MapInstance[activeCount];
}
ActiveMapInstances.Values.CopyTo(_activeMapInstancesBuffer, 0);
InstanceProcessor.UpdateInstanceControllers(_activeMapInstancesBuffer);
if (Options.Instance.Metrics.Enable)
{
MetricsRoot.Instance.Game.ActiveEntities.Record(globalEntities);
MetricsRoot.Instance.Game.ActiveEvents.Record(events);
MetricsRoot.Instance.Game.ProcessingEvents.Record(eventsProcessing);
MetricsRoot.Instance.Game.AutorunEvents.Record(autorunEvents);
MetricsRoot.Instance.Game.ActiveMaps.Record(ActiveMapInstances.Count);
MetricsRoot.Instance.Game.Players.Record(players);
MetricsRoot.Instance.Network.Clients.Record(Client.Instances?.Count ?? 0);
}
//End Resync of Active Maps
updateTimer = startTime + 250;
}
//Check our map update queues. If maps are ready to be updated based on our update intervals set in the server config then tell our thread pool to queue the map update as a work item.
lock (LogicLock)
{
if (Options.Instance.Processing.MapUpdateInterval != Options.Instance.Processing.ProjectileUpdateInterval)
{
while (MapInstanceProjectileQueue.TryPeek(out MapInstance result) && result.LastProjectileUpdateTime + Options.Instance.Processing.ProjectileUpdateInterval <= startTime)
{
if (MapInstanceProjectileQueue.TryDequeue(out MapInstance sameResult))
{
LogicPool.QueueWorkItem(UpdateMap, sameResult, true);
}
}
}
while (MapInstanceUpdateQueue.TryPeek(out MapInstance result) && result.LastRequestedUpdateTime + Options.Instance.Processing.MapUpdateInterval <= startTime)
{
if (MapInstanceUpdateQueue.TryDequeue(out MapInstance sameResult))
{
if (Options.Instance.Metrics.Enable)
{
var delay = Timing.Global.Milliseconds - (result.LastRequestedUpdateTime + Options.Instance.Processing.MapUpdateInterval);
MetricsRoot.Instance.Game.MapQueueUpdateOffset.Record(delay);
result.UpdateQueueStart = Timing.Global.Milliseconds;
}
LogicPool.QueueWorkItem(UpdateMap, sameResult, false);
}
}
}
Time.Update();
swCps++;
var endTime = Timing.Global.Milliseconds;
if (Timing.Global.Milliseconds > swCpsTimer)
{
_logicService.CyclesPerSecond = swCps;
swCps = 0;
var cyclesPerSecond = ApplicationContext.GetCurrentContext<IServerContext>().LogicService.CyclesPerSecond;
// Only call Console.Title when the value actually changed.
var newTitle = $"Intersect Server - CPS: {cyclesPerSecond}, Players: {players}, Active Maps: {ActiveMapInstances.Count}, Logic Threads: {LogicPool.ActiveThreads} ({LogicPool.InUseThreads} In Use), Pool Queue: {LogicPool.CurrentWorkItemsCount}, Idle: {LogicPool.IsIdle}";
if (newTitle != _lastConsoleTitle)
{
Console.Title = newTitle;
_lastConsoleTitle = newTitle;
}
if (Options.Instance.Metrics.Enable)
{
//Get Average CPU Usage for the last second
var currentCpuTime = Process.GetCurrentProcess().TotalProcessorTime;
var cpuUsedMs = (currentCpuTime - lastCpuTime).TotalMilliseconds;
var totalMsPassed = Timing.Global.Milliseconds - (swCpsTimer - 1000);
var cpuUsageTotal = (cpuUsedMs / (Environment.ProcessorCount * totalMsPassed)) * 100f;
lastCpuTime = currentCpuTime;
MetricsRoot.Instance.Application.Cpu.Record((int)cpuUsageTotal);
MetricsRoot.Instance.Application.Memory.Record(Process.GetCurrentProcess().PrivateMemorySize64);
MetricsRoot.Instance.Game.Cps.Record(cyclesPerSecond);
//Also Update Networking Metrics
MetricsRoot.Instance.Network.TotalBandwidth.Record(PacketHandler.ReceivedBytes + PacketSender.SentBytes);
MetricsRoot.Instance.Network.SentBytes.Record(PacketSender.SentBytes);
MetricsRoot.Instance.Network.SentPackets.Record(PacketSender.SentPackets);
MetricsRoot.Instance.Network.ReceivedBytes.Record(PacketHandler.ReceivedBytes);
MetricsRoot.Instance.Network.ReceivedPackets.Record(PacketHandler.ReceivedPackets);
MetricsRoot.Instance.Network.AcceptedBytes.Record(PacketHandler.AcceptedBytes);
MetricsRoot.Instance.Network.AcceptedPackets.Record(PacketHandler.AcceptedPackets);
MetricsRoot.Instance.Network.DroppedBytes.Record(PacketHandler.DroppedBytes);
MetricsRoot.Instance.Network.DroppedPackets.Record(PacketHandler.DroppedPackets);
PacketSender.ResetMetrics();
PacketHandler.ResetMetrics();
}
//Should we send out guild updates?
foreach (var guild in Guild.Guilds)
{
if (guild.Value.LastUpdateTime + Options.Instance.Guild.GuildUpdateInterval < Timing.Global.Milliseconds)
{
LogicPool.QueueWorkItem(guild.Value.Update);
}
}
swCpsTimer = Timing.Global.Milliseconds + 1000;
}
if (Options.Instance.Metrics.Enable)
{
//Record how our Thread Pools are Operating
MetricsRoot.Instance.Threading.LogicPoolActiveThreads.Record(LogicPool.ActiveThreads);
MetricsRoot.Instance.Threading.LogicPoolInUseThreads.Record(LogicPool.InUseThreads);
MetricsRoot.Instance.Threading.LogicPoolWorkItemsCount.Record(LogicPool.CurrentWorkItemsCount);
MetricsRoot.Instance.Threading.NetworkPoolActiveThreads.Record(NetworkPoolDataProvider.ActiveThreads);
MetricsRoot.Instance.Threading.NetworkPoolInUseThreads.Record(NetworkPoolDataProvider.InUseThreads);
MetricsRoot.Instance.Threading.NetworkPoolWorkItemsCount.Record(NetworkPoolDataProvider.CurrentWorkItemsCount);
MetricsRoot.Instance.Threading.DatabasePoolActiveThreads.Record(DbInterface.Pool.ActiveThreads);
MetricsRoot.Instance.Threading.DatabasePoolInUseThreads.Record(DbInterface.Pool.InUseThreads);
MetricsRoot.Instance.Threading.DatabasePoolWorkItemsCount.Record(DbInterface.Pool.CurrentWorkItemsCount);
ThreadPool.GetMaxThreads(out int maxWorkerThreads, out int maxIOThreads);
ThreadPool.GetAvailableThreads(out int availableWorkerThreads, out int availableIOThreads);
MetricsRoot.Instance.Threading.SystemPoolInUseWorkerThreads.Record(maxWorkerThreads - availableWorkerThreads);
MetricsRoot.Instance.Threading.SystemPoolInUseIOThreads.Record(maxIOThreads - availableIOThreads);
if (Timing.Global.Milliseconds > metricsTimer)
{
MetricsRoot.Instance.Capture();
foreach (var key in PacketSender.SentPacketTypes.Keys)
{
PacketSender.SentPacketTypes[key] = 0;
}
foreach (var key in PacketHandler.AcceptedPacketTypes.Keys)
{
PacketHandler.AcceptedPacketTypes[key] = 0;
}
metricsTimer = Timing.Global.Milliseconds + 5000;
}
}
if (saveServerVariablesTimer < endTime)
{
DbInterface.Pool.QueueWorkItem(DbInterface.SaveUpdatedServerVariables);
saveServerVariablesTimer = Timing.Global.Milliseconds + Options.Instance.Processing.DatabaseSaveServerVariablesInterval;
}
if (Options.Instance.Processing.CpsLock)
{
Thread.Sleep(1);
}
}
LogicPool.Shutdown();
}
catch (Exception exception)
{
ServerContext.DispatchUnhandledException(exception);
}
finally
{
ServerContext.Instance.RequestShutdown();
}
}
/// <summary>
/// Adds a map instance to the map update queues for our logic loop to start processing.
/// </summary>
/// <param name="mapInstance">The map instance in which to process in our queues.</param>
private void AddToQueue(MapInstance mapInstance)
{
if (Options.Instance.Processing.MapUpdateInterval != Options.Instance.Processing.ProjectileUpdateInterval)
{
MapInstanceProjectileQueue.Enqueue(mapInstance);
}
MapInstanceUpdateQueue.Enqueue(mapInstance);
ActiveMapInstances.Add(mapInstance.Id, mapInstance);
mapInstance.LastRequestedUpdateTime = Timing.Global.Milliseconds - Options.Instance.Processing.MapUpdateInterval;
}
/// <summary>
/// This function actually runs our map update function on the logic thread pool, and then re-queues our map for future updates if the map is still considered active.
/// </summary>
/// <param name="mapInstance">The <see cref="MapInstance"/> our thread updates.</param>
/// <param name="onlyProjectiles">If true only map projectiles are updated and not the entire map.</param>
private void UpdateMap(MapInstance mapInstance, bool onlyProjectiles)
{
try
{
if (onlyProjectiles)
{
mapInstance.UpdateProjectiles(Timing.Global.Milliseconds);
if (ActiveMapInstances.ContainsKey(mapInstance.Id))
{
MapInstanceProjectileQueue.Enqueue(mapInstance);
}
}
else
{
if (Options.Instance.Metrics.Enable)
{
var timeBeforeUpdate = Timing.Global.Milliseconds;
var desiredMapUpdateTime = mapInstance.LastRequestedUpdateTime + Options.Instance.Processing.MapUpdateInterval;
MetricsRoot.Instance.Game.MapUpdateQueuedTime.Record(timeBeforeUpdate - mapInstance.UpdateQueueStart);
mapInstance.Update(Timing.Global.Milliseconds);
var timeAfterUpdate = Timing.Global.Milliseconds;
MetricsRoot.Instance.Game.MapUpdateProcessingTime.Record(timeAfterUpdate - timeBeforeUpdate);
MetricsRoot.Instance.Game.MapTotalUpdateTime.Record(timeAfterUpdate - desiredMapUpdateTime);
}
else
{
mapInstance.Update(Timing.Global.Milliseconds);
}
if (ActiveMapInstances.ContainsKey(mapInstance.Id))
{
MapInstanceUpdateQueue.Enqueue(mapInstance);
}
}
}
catch (ThreadAbortException)
{
//Ignore if this pool is being shut down
}
catch (Exception exception)
{
ServerContext.DispatchUnhandledException(exception);
}
}
}
}