Small clean up for the job queue creation

This commit is contained in:
emeric
2026-06-02 17:22:33 +02:00
parent e73ea68785
commit 5b4eb8db8b
11 changed files with 21 additions and 16 deletions
@@ -26,12 +26,12 @@
namespace lms::scanner namespace lms::scanner
{ {
JobQueue::JobQueue(core::IJobScheduler& scheduler, std::size_t maxQueueSize, ProcessFunction processJobsDoneFunc, std::size_t batchSize, float _drainThreshold) JobQueue::JobQueue(core::IJobScheduler& scheduler, ProcessFunction processJobsDoneFunc, JobQueueParameters params)
: _scheduler{ scheduler } : _scheduler{ scheduler }
, _maxQueueSize{ maxQueueSize } , _maxQueueSize{ params.maxQueueSize }
, _processJobsDoneFunc{ std::move(processJobsDoneFunc) } , _processJobsDoneFunc{ std::move(processJobsDoneFunc) }
, _batchSize{ batchSize } , _batchSize{ params.processBatchSize }
, _drainThreshold{ _drainThreshold } , _drainThreshold{ params.drainThreshold }
{ {
assert(_scheduler.getJobsDoneCount() == 0); assert(_scheduler.getJobsDoneCount() == 0);
} }
@@ -32,14 +32,19 @@ namespace lms::core
namespace lms::scanner namespace lms::scanner
{ {
struct JobQueueParameters
{
std::size_t maxQueueSize = 20;
std::size_t processBatchSize = 1; // processBatchSize -> how many jobs done to notify at once using processJobsDoneFunc
float drainThreshold = 0.85F; // drainThreshold: fraction of maxQueueSize at which completed jobs are processed
};
class JobQueue class JobQueue
{ {
public: public:
using ProcessFunction = std::function<void(std::span<std::unique_ptr<core::IJob>>)>; using ProcessFunction = std::function<void(std::span<std::unique_ptr<core::IJob>>)>;
// processBatchSize -> how many jobs done to notify at once using processJobsDoneFunc JobQueue(core::IJobScheduler& scheduler, ProcessFunction processJobsDoneFunc, JobQueueParameters params = {});
// drainThreshold: fraction of maxQueueSize at which completed jobs are processed
JobQueue(core::IJobScheduler& scheduler, std::size_t maxQueueSize, ProcessFunction processJobsDoneFunc, std::size_t processBatchSize, float drainThreshold);
~JobQueue(); ~JobQueue();
JobQueue(const JobQueue&) = delete; JobQueue(const JobQueue&) = delete;
JobQueue& operator=(const JobQueue&) = delete; JobQueue& operator=(const JobQueue&) = delete;
@@ -381,7 +381,7 @@ namespace lms::scanner
}; };
{ {
JobQueue queue{ getJobScheduler(), 20, processJobsDone, 1, 0.85F }; JobQueue queue{ getJobScheduler(), processJobsDone };
db::ArtistId lastRetrievedArtistId{}; db::ArtistId lastRetrievedArtistId{};
db::IdRange<db::ArtistId> artistIdRange; db::IdRange<db::ArtistId> artistIdRange;
@@ -294,7 +294,7 @@ namespace lms::scanner
}; };
{ {
JobQueue queue{ getJobScheduler(), 20, processJobsDone, 1, 0.85F }; JobQueue queue{ getJobScheduler(), processJobsDone };
db::MediumId lastRetrievedMediumId{}; db::MediumId lastRetrievedMediumId{};
db::IdRange<db::MediumId> mediumIdRange; db::IdRange<db::MediumId> mediumIdRange;
@@ -183,7 +183,7 @@ namespace lms::scanner
}; };
{ {
JobQueue queue{ getJobScheduler(), 20, processJobsDone, 1, 0.85F }; JobQueue queue{ getJobScheduler(), processJobsDone };
db::PlayListFileId lastRetrievedId{}; db::PlayListFileId lastRetrievedId{};
db::IdRange<db::PlayListFileId> idRange; db::IdRange<db::PlayListFileId> idRange;
@@ -291,7 +291,7 @@ namespace lms::scanner
}; };
{ {
JobQueue queue{ getJobScheduler(), 20, processJobsDone, 1, 0.85F }; JobQueue queue{ getJobScheduler(), processJobsDone };
db::PlayListFileId lastPlayListFileId; db::PlayListFileId lastPlayListFileId;
db::IdRange<db::PlayListFileId> playListFileIdRange; db::IdRange<db::PlayListFileId> playListFileIdRange;
@@ -313,7 +313,7 @@ namespace lms::scanner
_progressCallback(context.currentStepStats); _progressCallback(context.currentStepStats);
}; };
JobQueue queue{ getJobScheduler(), 20, processJobsDone, 1, 0.85F }; JobQueue queue{ getJobScheduler(), processJobsDone };
db::ReleaseId lastRetrievedReleaseId{}; db::ReleaseId lastRetrievedReleaseId{};
db::IdRange<db::ReleaseId> artistIdRange; db::IdRange<db::ReleaseId> artistIdRange;
@@ -250,7 +250,7 @@ namespace lms::scanner
}; };
{ {
JobQueue queue{ getJobScheduler(), 20, processTracks, 1, 0.85F }; JobQueue queue{ getJobScheduler(), processTracks };
db::TrackId lastRetrievedTrackId; db::TrackId lastRetrievedTrackId;
db::IdRange<db::TrackId> trackIdRange; db::IdRange<db::TrackId> trackIdRange;
@@ -250,7 +250,7 @@ namespace lms::scanner
}; };
{ {
JobQueue queue{ getJobScheduler(), 50, processJobsDone, 1, 0.85F }; JobQueue queue{ getJobScheduler(), processJobsDone, { .maxQueueSize = 50 } };
ObjectIdType lastCheckedId; ObjectIdType lastCheckedId;
std::vector<FileToCheck<ObjectIdType>> filesToCheck; std::vector<FileToCheck<ObjectIdType>> filesToCheck;
@@ -214,7 +214,7 @@ namespace lms::scanner
} }; } };
{ {
JobQueue queue{ getJobScheduler(), 50, processResults, 1, 0.85F }; JobQueue queue{ getJobScheduler(), processResults, { .maxQueueSize = 50 } };
db::TrackId lastRetrievedTrackId; db::TrackId lastRetrievedTrackId;
TrackLocation trackLocation; TrackLocation trackLocation;
@@ -209,7 +209,7 @@ namespace lms::scanner
}; };
{ {
JobQueue queue{ getJobScheduler(), scanQueueMaxSize, processDoneJobs, processFileResultsBatchSize, drainRatio }; JobQueue queue{ getJobScheduler(), processDoneJobs, { .maxQueueSize = scanQueueMaxSize, .processBatchSize = processFileResultsBatchSize, .drainThreshold = drainRatio } };
std::vector<std::filesystem::directory_entry> filesToScan; std::vector<std::filesystem::directory_entry> filesToScan;