Reworked the queue items.

This commit is contained in:
2024-11-29 12:59:07 +01:00
parent 3034e66126
commit f2ace6f571
8 changed files with 152 additions and 91 deletions
+9 -15
View File
@@ -8,7 +8,8 @@ namespace Backend.Handler;
public class ContentFilter
{
private readonly ConcurrentQueue<QueueItem> _queue;
private readonly ConcurrentQueue<Filtered> _queue;
private readonly ConcurrentQueue<UnfilteredQueueItem> _unfilteredQueue;
private readonly DbHandler _dbHandler;
private readonly string _getDomainPort80;
private readonly string _getDomainPort443;
@@ -16,12 +17,13 @@ public class ContentFilter
private int _timeOut;
private readonly string _basePath;
public ContentFilter(ConcurrentQueue<QueueItem> queue, DbHandler dbHandler, string basePath)
public ContentFilter(ConcurrentQueue<Filtered> queue, ConcurrentQueue<UnfilteredQueueItem> unfilteredQueue, DbHandler dbHandler, string basePath)
{
_queue = queue;
_dbHandler = dbHandler;
_basePath = basePath;
_unfilteredQueue = unfilteredQueue;
_getDomainPort80 = $"{basePath}/Backend/Scripts/GetDomainNamePort80.sh";
_getDomainPort443 = $"{basePath}/Backend/Scripts/GetDomainNamePort443.sh";
@@ -63,14 +65,13 @@ public class ContentFilter
unfiltered.Filtered = true;
QueueItem superUnfilteredObject = new()
UnfilteredQueueItem superUnfilteredObject = new()
{
Unfiltered = unfiltered,
Operations = Operations.Update,
DbType = DbType.Unfiltered
Operations = Operations.Update
};
_queue.Enqueue(superUnfilteredObject);
_unfilteredQueue.Enqueue(superUnfilteredObject);
if (_dbHandler.FilteredIpExists(unfiltered.Ip))
{
@@ -82,14 +83,7 @@ public class ContentFilter
filtered.Port1 = unfiltered.Port1;
filtered.Port2 = unfiltered.Port2;
QueueItem superFilteredObject = new()
{
Filtered = filtered,
Operations = Operations.Insert,
DbType = DbType.Filtered
};
_queue.Enqueue(superFilteredObject);
_queue.Enqueue(filtered);
}
Thread.Sleep(_timeOut);
+12 -15
View File
@@ -17,18 +17,22 @@ public class ScanSettings
public class IpScanner
{
private readonly ConcurrentQueue<QueueItem> _queue;
private readonly ConcurrentQueue<Discarded> _discardedQueue;
private readonly ConcurrentQueue<UnfilteredQueueItem> _unfilteredQueue;
private readonly ConcurrentQueue<ScannerResumeObject> _resumeQueue;
private readonly DbHandler _dbHandler;
private bool _stop;
private int _timeout;
public IpScanner(ConcurrentQueue<QueueItem> queue, DbHandler dbHandler, ConcurrentQueue<Discarded> discardedQueue)
public IpScanner(ConcurrentQueue<UnfilteredQueueItem> unfilteredQueue, ConcurrentQueue<Discarded> discardedQueue,
ConcurrentQueue<ScannerResumeObject> resumeQueue, DbHandler dbHandler
)
{
_queue = queue;
_dbHandler = dbHandler;
_discardedQueue = discardedQueue;
_unfilteredQueue = unfilteredQueue;
_resumeQueue = resumeQueue;
SetTimeout(128);
}
@@ -167,7 +171,7 @@ public class IpScanner
continue;
}
_queue.Enqueue(CreateUnfilteredQueueItem(ip, ports));
_unfilteredQueue.Enqueue(CreateUnfilteredQueueItem(ip, ports));
}
if (_stop)
@@ -192,13 +196,7 @@ public class IpScanner
//Console.WriteLine($"Thread ({scanSettings.ThreadNumber}) is at index ({i}) out of ({scanSettings.End}). Remaining ({scanSettings.End - i})");
}
QueueItem resume = new()
{
ResumeObject = resumeObject,
Operations = Operations.Insert
};
_queue.Enqueue(resume);
_resumeQueue.Enqueue(resumeObject);
scanSettings.Handle!.Set();
}
@@ -212,7 +210,7 @@ public class IpScanner
};
}
private static QueueItem CreateUnfilteredQueueItem(Ip ip, (int, int) ports)
private static UnfilteredQueueItem CreateUnfilteredQueueItem(Ip ip, (int, int) ports)
{
Unfiltered unfiltered = new()
{
@@ -225,8 +223,7 @@ public class IpScanner
return new()
{
Unfiltered = unfiltered,
Operations = Operations.Insert,
DbType = DbType.Unfiltered
Operations = Operations.Insert
};
}
+29 -11
View File
@@ -17,33 +17,41 @@ public class ThreadHandler
public ThreadHandler(string path)
{
ConcurrentQueue<QueueItem> contentQueue = new();
ConcurrentQueue<Filtered> filteredQueue = new();
ConcurrentQueue<Discarded> discardedQueue = new();
ConcurrentQueue<UnfilteredQueueItem> unfilteredQueue = new();
ConcurrentQueue<ScannerResumeObject> scannerResumeQueue = new();
_dbHandler = new(contentQueue, discardedQueue, path);
_ipScanner = new(contentQueue, _dbHandler, discardedQueue);
_contentFilter = new(contentQueue, _dbHandler, path);
_dbHandler = new(filteredQueue, discardedQueue, unfilteredQueue, scannerResumeQueue, path);
_ipScanner = new(unfilteredQueue, discardedQueue, scannerResumeQueue, _dbHandler);
_contentFilter = new(filteredQueue, unfilteredQueue, _dbHandler, path);
_communication = new(_dbHandler, this, _ipScanner, _contentFilter, path);
}
public void Start()
{
//Thread scanner = new(StartScanner);
Thread scanner = new(StartScanner);
Thread indexer = new(StartContentFilter);
Thread database = new(StartDbHandler);
//Thread discarded = new(StartDiscardedDbHandler);
Thread discarded = new(StartDiscardedDbHandler);
Thread filtered = new(StartFilteredDbHandler);
Thread resume = new(StartResumeDbHandler);
Thread communication = new(StartCommunicationHandler);
//scanner.Start();
scanner.Start();
indexer.Start();
database.Start();
//discarded.Start();
discarded.Start();
filtered.Start();
resume.Start();
communication.Start();
//scanner.Join();
scanner.Join();
indexer.Join();
database.Join();
//discarded.Join();
discarded.Join();
filtered.Join();
resume.Join();
communication.Join();
}
@@ -73,7 +81,17 @@ public class ThreadHandler
private void StartDbHandler()
{
_dbHandler.StartContent();
_dbHandler.UnfilteredDbHandler();
}
private void StartFilteredDbHandler()
{
_dbHandler.FilteredDbHandler();
}
private void StartResumeDbHandler()
{
_dbHandler.ResumeDbHandler();
}
private void StartDiscardedDbHandler()