mirror of
https://github.com/elisspace/autopsy.git
synced 2026-09-06 02:24:30 +00:00
Fix IngestTasksScheduler, IngestManager concurrency issues
This commit is contained in:
@@ -49,12 +49,12 @@ final class IngestTasksScheduler {
|
||||
private static final int FAT_NTFS_FLAGS = TskData.TSK_FS_TYPE_ENUM.TSK_FS_TYPE_FAT12.getValue() | TskData.TSK_FS_TYPE_ENUM.TSK_FS_TYPE_FAT16.getValue() | TskData.TSK_FS_TYPE_ENUM.TSK_FS_TYPE_FAT32.getValue() | TskData.TSK_FS_TYPE_ENUM.TSK_FS_TYPE_NTFS.getValue();
|
||||
private static final Logger logger = Logger.getLogger(IngestTasksScheduler.class.getName());
|
||||
private static IngestTasksScheduler instance;
|
||||
private final List<DataSourceIngestTask> activeDataSourceTasks;
|
||||
private final DataSourceIngestTaskQueue dataSourceTaskQueue;
|
||||
private final TreeSet<FileIngestTask> rootFileTasks;
|
||||
private final Deque<FileIngestTask> directoryFileTasks;
|
||||
private final Deque<FileIngestTask> activeFileTasks;
|
||||
private final FileIngestTaskQueue fileTaskQueue;
|
||||
private final DataSourceIngestTaskQueue dataSourceTaskQueueForIngestThreads;
|
||||
private final List<DataSourceIngestTask> queuedAndRunningDataSourceTasks;
|
||||
private final TreeSet<FileIngestTask> rootFileTaskQueue;
|
||||
private final Deque<FileIngestTask> directoryFileTaskQueue;
|
||||
private final FileIngestTaskQueue fileTaskQueueForIngestThreads;
|
||||
private final List<FileIngestTask> queuedAndRunningFileTasks;
|
||||
|
||||
/**
|
||||
* Gets the ingest tasks scheduler singleton.
|
||||
@@ -70,41 +70,40 @@ final class IngestTasksScheduler {
|
||||
* Constructs an ingest tasks scheduler.
|
||||
*/
|
||||
private IngestTasksScheduler() {
|
||||
this.activeDataSourceTasks = new ArrayList<>();
|
||||
this.dataSourceTaskQueue = new DataSourceIngestTaskQueue();
|
||||
this.rootFileTasks = new TreeSet<>(new RootDirectoryTaskComparator());
|
||||
this.directoryFileTasks = new LinkedList<>();
|
||||
this.activeFileTasks = new LinkedList<>();
|
||||
this.fileTaskQueue = new FileIngestTaskQueue();
|
||||
this.queuedAndRunningDataSourceTasks = new LinkedList<>();
|
||||
this.dataSourceTaskQueueForIngestThreads = new DataSourceIngestTaskQueue();
|
||||
this.rootFileTaskQueue = new TreeSet<>(new RootDirectoryTaskComparator());
|
||||
this.directoryFileTaskQueue = new LinkedList<>();
|
||||
this.queuedAndRunningFileTasks = new LinkedList<>();
|
||||
this.fileTaskQueueForIngestThreads = new FileIngestTaskQueue();
|
||||
}
|
||||
|
||||
/**
|
||||
* Gets the data source level ingest tasks queue. The queue is a blocking
|
||||
* queue intended for use by data source ingest threads.
|
||||
* Gets the data source level ingest tasks queue. This queue is a blocking
|
||||
* queue intended for use by the ingest manager's data source ingest
|
||||
* threads.
|
||||
*
|
||||
* @return The queue.
|
||||
*/
|
||||
IngestTaskQueue getDataSourceIngestTaskQueue() {
|
||||
return this.dataSourceTaskQueue;
|
||||
return this.dataSourceTaskQueueForIngestThreads;
|
||||
}
|
||||
|
||||
/**
|
||||
* Gets the file level ingest tasks queue for file ingest threads. The queue
|
||||
* is a blocking queue intended for use by file ingest threads.
|
||||
* Gets the file level ingest tasks queue. This queue is a blocking queue
|
||||
* intended for use by the ingest manager's file ingest threads.
|
||||
*
|
||||
* @return The queue.
|
||||
*/
|
||||
IngestTaskQueue getFileIngestTaskQueue() {
|
||||
return this.fileTaskQueue;
|
||||
return this.fileTaskQueueForIngestThreads;
|
||||
}
|
||||
|
||||
/**
|
||||
* Schedules a data source level ingest task and file level ingest tasks for
|
||||
* a data source ingest job. Either all of the files in the data source or a
|
||||
* given subset of the files will be scheduled.
|
||||
* a data source ingest job.
|
||||
*
|
||||
* @param job The data source ingest job.
|
||||
* @param files A subset of the files for the data source, possibly empty.
|
||||
* @param job The data source ingest job.
|
||||
*/
|
||||
synchronized void scheduleIngestTasks(DataSourceIngestJob job) {
|
||||
if (!job.isCancelled()) {
|
||||
@@ -129,24 +128,21 @@ final class IngestTasksScheduler {
|
||||
synchronized void scheduleDataSourceIngestTask(DataSourceIngestJob job) {
|
||||
if (!job.isCancelled()) {
|
||||
DataSourceIngestTask task = new DataSourceIngestTask(job);
|
||||
this.activeDataSourceTasks.add(task);
|
||||
this.queuedAndRunningDataSourceTasks.add(task);
|
||||
try {
|
||||
this.dataSourceTaskQueue.add(task);
|
||||
this.dataSourceTaskQueueForIngestThreads.add(task);
|
||||
} catch (InterruptedException ex) {
|
||||
IngestTasksScheduler.logger.log(Level.INFO, "Ingest cancelled while data source ingest thread blocked on a full queue", ex);
|
||||
this.activeDataSourceTasks.remove(task);
|
||||
IngestTasksScheduler.logger.log(Level.INFO, "Ingest cancelled while a data source ingest thread was blocked on a full queue", ex);
|
||||
this.queuedAndRunningDataSourceTasks.remove(task);
|
||||
Thread.currentThread().interrupt();
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Schedules file level ingest tasks for a data source ingest job. Either
|
||||
* all of the files in the data source or a given subset of the files will
|
||||
* be scheduled.
|
||||
* Schedules file level ingest tasks for a data source ingest job.
|
||||
*
|
||||
* @param job The data source ingest job.
|
||||
* @param files A subset of the files for the data source, possibly empty.
|
||||
* @param job The data source ingest job.
|
||||
*/
|
||||
synchronized void scheduleFileIngestTasks(DataSourceIngestJob job) {
|
||||
if (!job.isCancelled()) {
|
||||
@@ -154,7 +150,7 @@ final class IngestTasksScheduler {
|
||||
for (AbstractFile file : candidateFiles) {
|
||||
FileIngestTask task = new FileIngestTask(job, file);
|
||||
if (IngestTasksScheduler.shouldEnqueueFileTask(task)) {
|
||||
this.rootFileTasks.add(task);
|
||||
this.rootFileTaskQueue.add(task);
|
||||
}
|
||||
}
|
||||
shuffleFileTaskQueues();
|
||||
@@ -162,7 +158,7 @@ final class IngestTasksScheduler {
|
||||
}
|
||||
|
||||
/**
|
||||
* Schedules file level ingest tasks for a subset of the files in a data
|
||||
* Schedules file level ingest tasks for a subset of the files for a data
|
||||
* source ingest job.
|
||||
*
|
||||
* @param job The data source ingest job.
|
||||
@@ -170,22 +166,27 @@ final class IngestTasksScheduler {
|
||||
*/
|
||||
synchronized void scheduleFileIngestTasks(DataSourceIngestJob job, Collection<AbstractFile> files) {
|
||||
if (!job.isCancelled()) {
|
||||
final Deque<FileIngestTask> newTasksForIngestThreads = new LinkedList<>();
|
||||
List<FileIngestTask> newTasks = new LinkedList<>();
|
||||
for (AbstractFile file : files) {
|
||||
/*
|
||||
* The file will be added directly to the front of the queue for
|
||||
* the ingest threads.
|
||||
* Put the file directly into the queue for the file ingest
|
||||
* threads, if it passes the filter for the job. The file is
|
||||
* added to the queue for the ingest threads BEFORE the other
|
||||
* queued tasks because the primary use case for this method is
|
||||
* adding derived files from a higher priority task that
|
||||
* preceded the tasks currently in the queue.
|
||||
*/
|
||||
FileIngestTask task = new FileIngestTask(job, file);
|
||||
if (shouldEnqueueFileTask(task)) {
|
||||
this.activeFileTasks.addLast(task); // RJCTODO: CHeck this in other method
|
||||
newTasksForIngestThreads.addLast(task);
|
||||
newTasks.add(task);
|
||||
}
|
||||
|
||||
/*
|
||||
* Add the children of the file, if any, either to the front of
|
||||
* the queue for the ingest threads, in front of the parent
|
||||
* directory task, or to the end directory task queue.
|
||||
* If the file or directory that was just queued has children,
|
||||
* try to queue tasks for the children. Each child task will go
|
||||
* into either the directory queue if it is a directory, or
|
||||
* directly into the queue for the file ingest threads, if it
|
||||
* passes the filter for the job.
|
||||
*/
|
||||
try {
|
||||
for (Content child : file.getChildren()) {
|
||||
@@ -193,10 +194,9 @@ final class IngestTasksScheduler {
|
||||
AbstractFile childFile = (AbstractFile) child;
|
||||
FileIngestTask childTask = new FileIngestTask(job, childFile);
|
||||
if (childFile.hasChildren()) {
|
||||
this.directoryFileTasks.add(childTask);
|
||||
this.directoryFileTaskQueue.add(childTask);
|
||||
} else if (shouldEnqueueFileTask(childTask)) {
|
||||
this.activeFileTasks.addLast(childTask);
|
||||
newTasksForIngestThreads.addFirst(childTask);
|
||||
newTasks.add(task);
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -206,15 +206,18 @@ final class IngestTasksScheduler {
|
||||
}
|
||||
|
||||
/*
|
||||
* Add the newly active tasks into the queue for the file ingest
|
||||
* threads, AFTER the higher priority tasks that are already queued.
|
||||
* The files are added to the queue for the ingest threads BEFORE
|
||||
* the other queued tasks because the primary use case for this
|
||||
* method is adding derived files from a higher priority task that
|
||||
* preceded the tasks currently in the queue.
|
||||
*/
|
||||
for (FileIngestTask newTask : newTasksForIngestThreads) {
|
||||
for (FileIngestTask newTask : newTasks) {
|
||||
try {
|
||||
this.fileTaskQueue.tasks.putLast(newTask);
|
||||
this.queuedAndRunningFileTasks.add(newTask);
|
||||
this.fileTaskQueueForIngestThreads.addFirst(newTask);
|
||||
} catch (InterruptedException ex) {
|
||||
IngestTasksScheduler.logger.log(Level.INFO, "Ingest cancelled while data source ingest thread blocked on a full queue", ex); // RJCTODO: Should this propagate? Correct message
|
||||
this.activeFileTasks.remove(newTask);
|
||||
this.queuedAndRunningFileTasks.remove(newTask);
|
||||
IngestTasksScheduler.logger.log(Level.INFO, "Ingest cancelled while blocked on a full file ingest threads queue", ex);
|
||||
Thread.currentThread().interrupt();
|
||||
break;
|
||||
}
|
||||
@@ -229,46 +232,45 @@ final class IngestTasksScheduler {
|
||||
* @param task The completed task.
|
||||
*/
|
||||
synchronized void notifyTaskCompleted(DataSourceIngestTask task) {
|
||||
this.activeDataSourceTasks.remove(task);
|
||||
shuffleFileTaskQueues();
|
||||
this.queuedAndRunningDataSourceTasks.remove(task);
|
||||
}
|
||||
|
||||
/**
|
||||
* Allows an ingest thread to notify this ingest task scheduler that a file
|
||||
* level task has been completed.
|
||||
*
|
||||
* @param task
|
||||
* @param task The completed task.
|
||||
*/
|
||||
synchronized void notifyTaskCompleted(FileIngestTask task) {
|
||||
this.activeFileTasks.remove(task);
|
||||
this.queuedAndRunningFileTasks.remove(task);
|
||||
shuffleFileTaskQueues();
|
||||
}
|
||||
|
||||
/**
|
||||
* Queries the task scheduler to determine whether or not all of the ingest
|
||||
* tasks for an ingest job have been completed.
|
||||
* tasks for a data source ingest job have been completed.
|
||||
*
|
||||
* @param job The data source ingest job
|
||||
* @param job The data source ingest job.
|
||||
*
|
||||
* @return True or false.
|
||||
*/
|
||||
synchronized boolean tasksForJobAreCompleted(DataSourceIngestJob job) {
|
||||
return !hasTasksForJob(this.activeDataSourceTasks, job)
|
||||
&& !hasTasksForJob(this.rootFileTasks, job)
|
||||
&& !hasTasksForJob(this.directoryFileTasks, job)
|
||||
&& !hasTasksForJob(this.activeFileTasks, job);
|
||||
return !hasTasksForJob(this.queuedAndRunningDataSourceTasks, job)
|
||||
&& !hasTasksForJob(this.rootFileTaskQueue, job)
|
||||
&& !hasTasksForJob(this.directoryFileTaskQueue, job)
|
||||
&& !hasTasksForJob(this.queuedAndRunningFileTasks, job);
|
||||
}
|
||||
|
||||
/**
|
||||
* Clears the "upstream" task scheduling queues for a data source ingest
|
||||
* job, but does nothing about tasks that have already been activated, i.e.,
|
||||
* moved into the queue that is consumed by the file ingest threads.
|
||||
* job, but does nothing about tasks that have already been moved into the
|
||||
* queue that is consumed by the file ingest threads.
|
||||
*
|
||||
* @param job The data source ingest job
|
||||
* @param job The data source ingest job.
|
||||
*/
|
||||
synchronized void cancelPendingTasksForIngestJob(DataSourceIngestJob job) {
|
||||
this.removeTasksForJob(this.rootFileTasks, job);
|
||||
this.removeTasksForJob(this.directoryFileTasks, job);
|
||||
this.removeTasksForJob(this.rootFileTaskQueue, job);
|
||||
this.removeTasksForJob(this.directoryFileTaskQueue, job);
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -276,8 +278,9 @@ final class IngestTasksScheduler {
|
||||
* files and virtual directories for a data source. Used to create file
|
||||
* tasks to put into the root directories queue.
|
||||
*
|
||||
* @param dataSource The data source.
|
||||
* @param topLevelFiles The top level files are added to this list.
|
||||
* @param dataSource The data source.
|
||||
*
|
||||
* @return The top level files.
|
||||
*/
|
||||
private static List<AbstractFile> getTopLevelFiles(Content dataSource) {
|
||||
List<AbstractFile> topLevelFiles = new ArrayList<>();
|
||||
@@ -312,42 +315,53 @@ final class IngestTasksScheduler {
|
||||
}
|
||||
|
||||
/**
|
||||
* Intelligently queues file ingest tasks for the ingest manager's file
|
||||
* ingest threads by "shuffling" them through a sequence of queues that
|
||||
* allows for the interleaving of tasks from different data source ingest
|
||||
* jobs based on priority. The sequence of queues is:
|
||||
* Schedules file ingest tasks for the ingest manager's file ingest threads
|
||||
* by "shuffling" them through a sequence of three queues that allows for
|
||||
* the interleaving of tasks from different data source ingest jobs based on
|
||||
* priority. The sequence of queues is:
|
||||
*
|
||||
* 1. The root file tasks priority queue, which contains the tasks for the
|
||||
* roots of data source content sub trees that are being analyzed for the
|
||||
* data source ingest jobs. For example, typical root tasks for a disk image
|
||||
* data source would be the tasks for the contents of the root directories
|
||||
* of the file systems. This queue is a priority queue that attempts to
|
||||
* ensure that user directory content is analyzed before general file system
|
||||
* content. It feeds into the directory tasks queue.
|
||||
* 1. The root file tasks priority queue, which contains file tasks for the
|
||||
* root objects of the data sources that are being analyzed. For example,
|
||||
* the root tasks for a disk image data source are typically the tasks for
|
||||
* the contents of the root directories of the file systems. This queue is a
|
||||
* priority queue that attempts to ensure that user directory content is
|
||||
* analyzed before general file system content. It feeds into the directory
|
||||
* tasks queue.
|
||||
*
|
||||
* 2. The directory file tasks queue, which contains directory tasks
|
||||
* discovered in the descent through the content sub trees that are being
|
||||
* analyzed for the data source ingest jobs. It feeds into the active tasks
|
||||
* queue.
|
||||
* 2. The directory file tasks queue, which contains root file tasks
|
||||
* shuffled out of the root tasks queue, plus directory tasks discovered in
|
||||
* the descent from the root tasks to the final leaf tasks in the content
|
||||
* trees that are being analyzed for the data source ingest jobs. This queue
|
||||
* is a FIFO queue. It feeds into the file tasks queue for the ingest
|
||||
* manager's file ingest threads.
|
||||
*
|
||||
* 3. The active file tasks queue, a queue of the file tasks that are either
|
||||
* in the tasks queue for the file ingest threads or are in the process of
|
||||
* being analyzed in a file ingest thread.
|
||||
* 3. The file tasks queue for the ingest manager's file ingest threads.
|
||||
* This queue is a blocking deque that is FIFO during a shuffle to maintain
|
||||
* task prioritization, but LIFO when adding derived files to it directly
|
||||
* during ingest. The reason for the LIFO additions is to give priority
|
||||
* derived files of priority files.
|
||||
*
|
||||
* 4. The file tasks queue for the ingest manager's file ingest threads.
|
||||
* There is a fourth collection of file tasks, a "tracking" list, that keeps
|
||||
* track of the file tasks that are either in the tasks queue for the file
|
||||
* ingest threads, or are in the process of being analyzed in a file ingest
|
||||
* thread. This queue is vital to the ingest task scheduler's ability to
|
||||
* determine when all of the ingest tasks for a data source ingest job have
|
||||
* been completed. It is also used to drive this shuffling algorithm -
|
||||
* whenever this list is empty, the two "upstream" queues are "shuffled" to
|
||||
* queue more tasks for the file ingest threads.
|
||||
*/
|
||||
synchronized private void shuffleFileTaskQueues() {
|
||||
final Deque<FileIngestTask> newTasksForIngestThreads = new LinkedList<>();
|
||||
while (this.activeFileTasks.isEmpty()) {
|
||||
List<FileIngestTask> newTasks = new LinkedList<>();
|
||||
while (this.queuedAndRunningFileTasks.isEmpty()) {
|
||||
/*
|
||||
* If the directory file task queue is empty, move the highest
|
||||
* priority root file task, if there is one, into it. If both the
|
||||
* root and the directory file task queuess are empty, there is
|
||||
* nothing left to do.
|
||||
* root and the directory file task queues are empty, there is
|
||||
* nothing left to shuffle, so exit.
|
||||
*/
|
||||
if (this.directoryFileTasks.isEmpty()) {
|
||||
if (!this.rootFileTasks.isEmpty()) {
|
||||
this.directoryFileTasks.add(this.rootFileTasks.pollFirst());
|
||||
if (this.directoryFileTaskQueue.isEmpty()) {
|
||||
if (!this.rootFileTaskQueue.isEmpty()) {
|
||||
this.directoryFileTaskQueue.add(this.rootFileTaskQueue.pollFirst());
|
||||
} else {
|
||||
return;
|
||||
}
|
||||
@@ -355,21 +369,22 @@ final class IngestTasksScheduler {
|
||||
|
||||
/*
|
||||
* Try to move the next task from the directory task queue into the
|
||||
* active file tasks tracking queue, if it passes the filter for the
|
||||
* job.
|
||||
* queue for the file ingest threads, if it passes the filter for
|
||||
* the job. The file is added to the queue for the ingest threads
|
||||
* AFTER the higher priority tasks that preceded it.
|
||||
*/
|
||||
final FileIngestTask directoryTask = this.directoryFileTasks.pollLast();
|
||||
final FileIngestTask directoryTask = this.directoryFileTaskQueue.pollLast();
|
||||
if (shouldEnqueueFileTask(directoryTask)) {
|
||||
this.activeFileTasks.addFirst(directoryTask);
|
||||
newTasksForIngestThreads.addFirst(directoryTask);
|
||||
newTasks.add(directoryTask);
|
||||
}
|
||||
|
||||
/*
|
||||
* If the file or directory from the next that was just activated
|
||||
* has children, try to queue tasks for the children. Each child
|
||||
* will go into the directory task queue if it is a directory, or
|
||||
* into the active file tasks tracking queue if it passes the filter
|
||||
* for the job.
|
||||
* If the directory (or root level file) that was just queued has
|
||||
* children, try to queue tasks for the children. Each child task
|
||||
* will go into either the directory queue if it is a directory, or
|
||||
* into the queue for the file ingest threads, if it passes the
|
||||
* filter for the job. The file is added to the queue for the ingest
|
||||
* threads AFTER the higher priority tasks that preceded it.
|
||||
*/
|
||||
final AbstractFile directory = directoryTask.getFile();
|
||||
try {
|
||||
@@ -378,38 +393,30 @@ final class IngestTasksScheduler {
|
||||
AbstractFile childFile = (AbstractFile) child;
|
||||
FileIngestTask childTask = new FileIngestTask(directoryTask.getIngestJob(), childFile);
|
||||
if (childFile.hasChildren()) {
|
||||
this.directoryFileTasks.add(childTask);
|
||||
this.directoryFileTaskQueue.add(childTask);
|
||||
} else if (shouldEnqueueFileTask(childTask)) {
|
||||
this.activeFileTasks.add(directoryTask);
|
||||
/*
|
||||
* Queue the child file tasks for the ingest threads
|
||||
* in front of parent directory tasks.
|
||||
*/
|
||||
newTasksForIngestThreads.addFirst(directoryTask);
|
||||
newTasks.add(childTask);
|
||||
}
|
||||
}
|
||||
}
|
||||
} catch (TskCoreException ex) {
|
||||
logger.log(Level.SEVERE, String.format("Error getting the children of %s (objId=%d)", directory.getName(), directory.getId()), ex); //NON-NLS
|
||||
}
|
||||
}
|
||||
|
||||
/*
|
||||
* Add the newly active tasks into the queue for the file ingest
|
||||
* threads, AFTER the higher priority tasks that are already queued.
|
||||
*/
|
||||
for (FileIngestTask newTask : newTasksForIngestThreads) {
|
||||
try {
|
||||
this.fileTaskQueue.tasks.putLast(newTask);
|
||||
} catch (InterruptedException ex) {
|
||||
/**
|
||||
* The current thread was interrupted while blocked on a
|
||||
* full queue. Discard the task and reset the interrupted
|
||||
* flag.
|
||||
*/
|
||||
IngestTasksScheduler.logger.log(Level.INFO, "Ingest cancelled while data source ingest thread blocked on a full queue", ex); // RJCTODO: Should this propagate? Correct message
|
||||
this.activeFileTasks.remove(newTask);
|
||||
Thread.currentThread().interrupt();
|
||||
}
|
||||
/*
|
||||
* The files are added to the queue for the ingest threads AFTER the
|
||||
* higher priority tasks that preceded them.
|
||||
*/
|
||||
for (FileIngestTask newTask : newTasks) {
|
||||
try {
|
||||
this.queuedAndRunningFileTasks.add(newTask);
|
||||
this.fileTaskQueueForIngestThreads.addFirst(newTask);
|
||||
} catch (InterruptedException ex) {
|
||||
this.queuedAndRunningFileTasks.remove(newTask);
|
||||
IngestTasksScheduler.logger.log(Level.INFO, "Ingest cancelled while blocked on a full file ingest threads queue", ex);
|
||||
Thread.currentThread().interrupt();
|
||||
break;
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -674,32 +681,40 @@ final class IngestTasksScheduler {
|
||||
}
|
||||
|
||||
/**
|
||||
* Wraps access to pending data source ingest tasks in the interface
|
||||
* required by the ingest threads.
|
||||
* A blocking queue of data source ingest tasks for the ingest manager's
|
||||
* data source ingest threads.
|
||||
*/
|
||||
private final class DataSourceIngestTaskQueue implements IngestTaskQueue {
|
||||
|
||||
private final BlockingQueue<DataSourceIngestTask> tasks = new LinkedBlockingQueue<>();
|
||||
|
||||
private void add(DataSourceIngestTask task) throws InterruptedException {
|
||||
this.tasks.put(task);
|
||||
}
|
||||
|
||||
@Override
|
||||
public IngestTask getNextTask() throws InterruptedException {
|
||||
return tasks.take();
|
||||
}
|
||||
|
||||
private void add(DataSourceIngestTask task) throws InterruptedException {
|
||||
this.tasks.put(task);
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
/**
|
||||
* Wraps access to pending file ingest tasks in the interface required by
|
||||
* the ingest threads.
|
||||
* A blocking, LIFO queue of data source ingest tasks for the ingest
|
||||
* manager's data source ingest threads.
|
||||
*/
|
||||
private final class FileIngestTaskQueue implements IngestTaskQueue {
|
||||
|
||||
private final BlockingDeque<FileIngestTask> tasks = new LinkedBlockingDeque<>();
|
||||
|
||||
private void addFirst(FileIngestTask task) throws InterruptedException {
|
||||
this.tasks.putFirst(task);
|
||||
}
|
||||
|
||||
private void addLast(FileIngestTask task) throws InterruptedException {
|
||||
this.tasks.putLast(task);
|
||||
}
|
||||
|
||||
@Override
|
||||
public IngestTask getNextTask() throws InterruptedException {
|
||||
return tasks.takeFirst();
|
||||
@@ -726,11 +741,11 @@ final class IngestTasksScheduler {
|
||||
*/
|
||||
IngestJobTasksSnapshot(long jobId) {
|
||||
this.jobId = jobId;
|
||||
this.rootQueueSize = countTasksForJob(IngestTasksScheduler.this.rootFileTasks, jobId);
|
||||
this.dirQueueSize = countTasksForJob(IngestTasksScheduler.this.directoryFileTasks, jobId);
|
||||
this.fileQueueSize = countTasksForJob(IngestTasksScheduler.this.fileTaskQueue.tasks, jobId);
|
||||
this.dsQueueSize = countTasksForJob(IngestTasksScheduler.this.dataSourceTaskQueue.tasks, jobId);
|
||||
this.runningListSize = countTasksForJob(IngestTasksScheduler.this.activeDataSourceTasks, jobId) + countTasksForJob(IngestTasksScheduler.this.activeFileTasks, jobId);
|
||||
this.rootQueueSize = countTasksForJob(IngestTasksScheduler.this.rootFileTaskQueue, jobId);
|
||||
this.dirQueueSize = countTasksForJob(IngestTasksScheduler.this.directoryFileTaskQueue, jobId);
|
||||
this.fileQueueSize = countTasksForJob(IngestTasksScheduler.this.fileTaskQueueForIngestThreads.tasks, jobId);
|
||||
this.dsQueueSize = countTasksForJob(IngestTasksScheduler.this.dataSourceTaskQueueForIngestThreads.tasks, jobId);
|
||||
this.runningListSize = countTasksForJob(IngestTasksScheduler.this.queuedAndRunningDataSourceTasks, jobId) + countTasksForJob(IngestTasksScheduler.this.queuedAndRunningFileTasks, jobId);
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
Reference in New Issue
Block a user