Fixed a new memory leak
This commit is contained in:
@@ -1,7 +1,5 @@
|
||||
using System.Collections.Concurrent;
|
||||
using System.Diagnostics;
|
||||
using Backend.Helper;
|
||||
using Models.Handler;
|
||||
using Models.Model.Backend;
|
||||
|
||||
namespace Backend.Handler;
|
||||
@@ -10,23 +8,20 @@ public class IpFilterHandler
|
||||
{
|
||||
private readonly ConcurrentQueue<Discarded> _discardedQueue;
|
||||
private readonly ConcurrentQueue<UnfilteredQueueItem> _unfilteredQueue;
|
||||
private readonly ConcurrentQueue<FilterQueueItem> _preFilteredQueue;
|
||||
private DbHandler _dbHandler;
|
||||
private readonly ConcurrentQueue<Ip> _preFilteredQueue;
|
||||
private ThreadHandler _threadHandler;
|
||||
private bool _stop;
|
||||
private bool _fillerStop;
|
||||
private bool _done;
|
||||
private bool _stopAutoscaledThreads;
|
||||
private int _timeout;
|
||||
|
||||
public IpFilterHandler(ConcurrentQueue<Discarded> discardedQueue,
|
||||
ConcurrentQueue<UnfilteredQueueItem> unfilteredQueue,
|
||||
ConcurrentQueue<FilterQueueItem> preFilteredQueue,
|
||||
DbHandler dbHandler, ThreadHandler threadHandler)
|
||||
ConcurrentQueue<Ip> preFilteredQueue, ThreadHandler threadHandler)
|
||||
{
|
||||
_discardedQueue = discardedQueue;
|
||||
_unfilteredQueue = unfilteredQueue;
|
||||
_preFilteredQueue = preFilteredQueue;
|
||||
_dbHandler = dbHandler;
|
||||
_threadHandler = threadHandler;
|
||||
|
||||
_timeout = 16;
|
||||
@@ -67,64 +62,15 @@ public class IpFilterHandler
|
||||
|
||||
return waitHandles;
|
||||
}
|
||||
|
||||
public void AutoScaler()
|
||||
{
|
||||
int i = 0;
|
||||
int j = 0;
|
||||
|
||||
while (!_stop)
|
||||
{
|
||||
if (_preFilteredQueue.Count >= 2000)
|
||||
{
|
||||
if (i == 10)
|
||||
{
|
||||
_stopAutoscaledThreads = false;
|
||||
Console.WriteLine("Autoscaler started");
|
||||
|
||||
while (!_stopAutoscaledThreads)
|
||||
{
|
||||
if (_preFilteredQueue.Count <= 2000)
|
||||
{
|
||||
if (j == 1000)
|
||||
{
|
||||
_stopAutoscaledThreads = true;
|
||||
}
|
||||
|
||||
j++;
|
||||
|
||||
Thread.Sleep(128);
|
||||
}
|
||||
else
|
||||
{
|
||||
EventWaitHandle handle = new(false, EventResetMode.ManualReset);
|
||||
Thread f = new (Filter_AutoScaler!);
|
||||
f.Start(handle);
|
||||
Thread.Sleep(16);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
i++;
|
||||
}
|
||||
else
|
||||
{
|
||||
i = 0;
|
||||
j = 0;
|
||||
}
|
||||
|
||||
Thread.Sleep(128);
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
private void Filter(object obj)
|
||||
{
|
||||
int counter = 0;
|
||||
while (!_stop)
|
||||
{
|
||||
if (_preFilteredQueue.IsEmpty && _fillerStop)
|
||||
if (_preFilteredQueue.IsEmpty)
|
||||
{
|
||||
if (counter == 100)
|
||||
if (counter == 30_000)
|
||||
{
|
||||
_threadHandler.Stop();
|
||||
_stop = true;
|
||||
@@ -132,63 +78,21 @@ public class IpFilterHandler
|
||||
|
||||
counter++;
|
||||
Thread.Sleep(128);
|
||||
}
|
||||
|
||||
_preFilteredQueue.TryDequeue(out FilterQueueItem item);
|
||||
|
||||
(int, int) ports = TcpClientHelper.CheckPort(item.Ip, 80, 443);
|
||||
|
||||
if (ports is { Item1: 0, Item2: 0 })
|
||||
{
|
||||
_discardedQueue.Enqueue(CreateDiscardedQueueItem(item.Ip, item.ResponseCode));
|
||||
continue;
|
||||
}
|
||||
|
||||
_unfilteredQueue.Enqueue(CreateUnfilteredQueueItem(item.Ip, ports));
|
||||
}
|
||||
|
||||
((EventWaitHandle) obj).Set();
|
||||
}
|
||||
|
||||
public void FillFilterQueue()
|
||||
{
|
||||
Console.WriteLine("Fill FilterQueue started.");
|
||||
while (!_stop)
|
||||
{
|
||||
if (_preFilteredQueue.Count > 500) continue;
|
||||
|
||||
if (_dbHandler.GetPreFilterQueueItem(out FilterQueueItem item))
|
||||
{
|
||||
_preFilteredQueue.Enqueue(item);
|
||||
}
|
||||
else
|
||||
{
|
||||
_fillerStop = true;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private void Filter_AutoScaler(object obj)
|
||||
{
|
||||
while (!_stopAutoscaledThreads)
|
||||
{
|
||||
if (_preFilteredQueue.IsEmpty)
|
||||
{
|
||||
Thread.Sleep(_timeout);
|
||||
|
||||
continue;
|
||||
}
|
||||
|
||||
_preFilteredQueue.TryDequeue(out FilterQueueItem item);
|
||||
_preFilteredQueue.TryDequeue(out Ip item);
|
||||
|
||||
(int, int) ports = TcpClientHelper.CheckPort(item.Ip, 80, 443);
|
||||
(int, int) ports = TcpClientHelper.CheckPort(item, 80, 443);
|
||||
|
||||
if (ports is { Item1: 0, Item2: 0 })
|
||||
{
|
||||
_discardedQueue.Enqueue(CreateDiscardedQueueItem(item.Ip, item.ResponseCode));
|
||||
_discardedQueue.Enqueue(CreateDiscardedQueueItem(item, item.ResponseCode));
|
||||
continue;
|
||||
}
|
||||
|
||||
_unfilteredQueue.Enqueue(CreateUnfilteredQueueItem(item.Ip, ports));
|
||||
_unfilteredQueue.Enqueue(CreateUnfilteredQueueItem(item, ports));
|
||||
}
|
||||
|
||||
((EventWaitHandle) obj).Set();
|
||||
|
||||
@@ -22,7 +22,7 @@ public class ScanSettings
|
||||
public class IpScanner
|
||||
{
|
||||
private readonly ConcurrentQueue<Discarded> _discardedQueue;
|
||||
private readonly ConcurrentQueue<FilterQueueItem> _preFilteredQueue;
|
||||
private readonly ConcurrentQueue<Ip> _preFilteredQueue;
|
||||
private readonly ConcurrentQueue<ScannerResumeObject> _resumeQueue;
|
||||
private readonly DbHandler _dbHandler;
|
||||
private bool _stop;
|
||||
@@ -30,7 +30,7 @@ public class IpScanner
|
||||
|
||||
public IpScanner(ConcurrentQueue<Discarded> discardedQueue,
|
||||
ConcurrentQueue<ScannerResumeObject> resumeQueue, DbHandler dbHandler,
|
||||
ConcurrentQueue<FilterQueueItem> preFilteredQueue)
|
||||
ConcurrentQueue<Ip> preFilteredQueue)
|
||||
{
|
||||
_dbHandler = dbHandler;
|
||||
_preFilteredQueue = preFilteredQueue;
|
||||
@@ -283,15 +283,11 @@ public class IpScanner
|
||||
};
|
||||
}
|
||||
|
||||
private static FilterQueueItem CreateUnfilteredQueueItem(Ip ip, int responseCode)
|
||||
private static Ip CreateUnfilteredQueueItem(Ip ip, int responseCode)
|
||||
{
|
||||
FilterQueueItem filterQueueItem = new()
|
||||
{
|
||||
Ip = ip,
|
||||
ResponseCode = responseCode
|
||||
};
|
||||
ip.ResponseCode = responseCode;
|
||||
|
||||
return filterQueueItem;
|
||||
return ip;
|
||||
}
|
||||
|
||||
public void Stop()
|
||||
|
||||
@@ -15,22 +15,22 @@ public class ThreadHandler
|
||||
private bool _contentFilterStopped;
|
||||
private bool _ipFilterStopped;
|
||||
|
||||
private bool _stage1 = true;
|
||||
private bool _stage2;
|
||||
private bool _stage3;
|
||||
private bool _stage1 = false;
|
||||
private bool _stage2 = true;
|
||||
private bool _stage3 = false;
|
||||
|
||||
ConcurrentQueue<Filtered> filteredQueue = new();
|
||||
ConcurrentQueue<Discarded> discardedQueue = new();
|
||||
ConcurrentQueue<UnfilteredQueueItem> unfilteredQueue = new();
|
||||
ConcurrentQueue<ScannerResumeObject> scannerResumeQueue = new();
|
||||
ConcurrentQueue<FilterQueueItem> preFilteredQueue = new();
|
||||
ConcurrentQueue<Ip> preFilteredQueue = new();
|
||||
|
||||
public ThreadHandler(string path)
|
||||
{
|
||||
_dbHandler = new(filteredQueue, discardedQueue, unfilteredQueue, scannerResumeQueue, preFilteredQueue, path);
|
||||
_ipScanner = new(discardedQueue, scannerResumeQueue, _dbHandler, preFilteredQueue);
|
||||
_contentFilter = new(filteredQueue, unfilteredQueue, _dbHandler, path, this);
|
||||
_ipFilterHandler = new(discardedQueue, unfilteredQueue, preFilteredQueue, _dbHandler, this);
|
||||
_ipFilterHandler = new(discardedQueue, unfilteredQueue, preFilteredQueue, this);
|
||||
}
|
||||
|
||||
public void Start()
|
||||
@@ -42,19 +42,20 @@ public class ThreadHandler
|
||||
Thread discarded = new(StartDiscardedDbHandler);
|
||||
Thread filtered = new(StartFilteredDbHandler);
|
||||
Thread resume = new(StartResumeDbHandler);
|
||||
Thread ipFilterAutoScaler = new(StartIpFilterAutoScaler);
|
||||
Thread contentFilterThread = new(StartContentFilterThread);
|
||||
Thread prefilterDb = new(StartPreFilterDbHandler);
|
||||
Thread fillIpFilterQueue = new(StartFillIpFilterQueue);
|
||||
//Thread check = new(CheckQueue);
|
||||
|
||||
|
||||
//check.Start();
|
||||
|
||||
if (_stage1)
|
||||
{
|
||||
discarded.Start(); // de-queues from discardedQueue
|
||||
prefilterDb.Start(); // de-queues from preFilteredQueue
|
||||
scanner.Start(); // en-queues to discardedQueue and preFilteredQueue
|
||||
resume.Start(); // de-queues from resumeQueue
|
||||
|
||||
|
||||
scanner.Join();
|
||||
Stop();
|
||||
|
||||
@@ -62,22 +63,21 @@ public class ThreadHandler
|
||||
prefilterDb.Join();
|
||||
resume.Join();
|
||||
}
|
||||
|
||||
|
||||
if (_stage2)
|
||||
{
|
||||
database.Start(); // de-queues from unfilteredQueue
|
||||
discarded.Start(); // de-queues from discardedQueue
|
||||
ipFilter.Start(); // en-queues to discardedQueue and unfilteredQueue
|
||||
ipFilterAutoScaler.Start(); // de-queues from preFilteredQueue, en-queues to discardedQueue and unfilteredQueue
|
||||
fillIpFilterQueue.Start(); // reads from preFiltered database, en-queues to preFilteredQueue
|
||||
|
||||
|
||||
ipFilter.Join();
|
||||
Stop();
|
||||
database.Join();
|
||||
discarded.Join();
|
||||
ipFilter.Join();
|
||||
ipFilterAutoScaler.Join();
|
||||
fillIpFilterQueue.Join();
|
||||
}
|
||||
|
||||
|
||||
if (_stage3)
|
||||
{
|
||||
filtered.Start(); // de-queues from filteredQueue
|
||||
@@ -91,7 +91,7 @@ public class ThreadHandler
|
||||
indexer.Join();
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
private void CheckQueue()
|
||||
{
|
||||
while (true)
|
||||
@@ -105,13 +105,14 @@ public class ThreadHandler
|
||||
Thread.Sleep(5);
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
private void StartScanner()
|
||||
{
|
||||
Thread.Sleep(15000); // Let the database handler instantiate and warm up first.
|
||||
|
||||
List<WaitHandle[]> wait = _ipScanner.Start(256);
|
||||
|
||||
|
||||
for (int i = 0; i < wait.Count; i++)
|
||||
{
|
||||
WaitHandle.WaitAll(wait[i]);
|
||||
@@ -135,26 +136,21 @@ public class ThreadHandler
|
||||
|
||||
private void StartContentFilterThread()
|
||||
{
|
||||
WaitHandle[] wait = _contentFilter.StartFilterThread(64);
|
||||
WaitHandle[] wait = _contentFilter.StartFilterThread(8);
|
||||
|
||||
WaitHandle.WaitAll(wait);
|
||||
}
|
||||
|
||||
private void StartIpFilterAutoScaler()
|
||||
{
|
||||
_ipFilterHandler.AutoScaler();
|
||||
}
|
||||
|
||||
private void StartFillIpFilterQueue()
|
||||
{
|
||||
_ipFilterHandler.FillFilterQueue();
|
||||
_dbHandler.GetPreFilterQueueItem();
|
||||
}
|
||||
|
||||
private void StartIpFilter()
|
||||
{
|
||||
Thread.Sleep(1000);
|
||||
|
||||
List<WaitHandle[]> wait = _ipFilterHandler.Start(256);
|
||||
List<WaitHandle[]> wait = _ipFilterHandler.Start(1024);
|
||||
|
||||
for (int i = 0; i < wait.Count; i++)
|
||||
{
|
||||
|
||||
Reference in New Issue
Block a user