From 0d4a2315cd67659f48eed00a126b7ba64d2ba9ce Mon Sep 17 00:00:00 2001 From: Richard Cordovano Date: Mon, 19 Mar 2018 19:09:36 -0400 Subject: [PATCH] Fix IngestTasksScheduler, IngestManager concurrency issues --- .../autopsy/ingest/IngestTasksScheduler.java | 295 +++++++++--------- 1 file changed, 155 insertions(+), 140 deletions(-) diff --git a/Core/src/org/sleuthkit/autopsy/ingest/IngestTasksScheduler.java b/Core/src/org/sleuthkit/autopsy/ingest/IngestTasksScheduler.java index 7fea572eed..2bc1769bc6 100644 --- a/Core/src/org/sleuthkit/autopsy/ingest/IngestTasksScheduler.java +++ b/Core/src/org/sleuthkit/autopsy/ingest/IngestTasksScheduler.java @@ -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 activeDataSourceTasks; - private final DataSourceIngestTaskQueue dataSourceTaskQueue; - private final TreeSet rootFileTasks; - private final Deque directoryFileTasks; - private final Deque activeFileTasks; - private final FileIngestTaskQueue fileTaskQueue; + private final DataSourceIngestTaskQueue dataSourceTaskQueueForIngestThreads; + private final List queuedAndRunningDataSourceTasks; + private final TreeSet rootFileTaskQueue; + private final Deque directoryFileTaskQueue; + private final FileIngestTaskQueue fileTaskQueueForIngestThreads; + private final List 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 files) { if (!job.isCancelled()) { - final Deque newTasksForIngestThreads = new LinkedList<>(); + List 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 getTopLevelFiles(Content dataSource) { List 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 newTasksForIngestThreads = new LinkedList<>(); - while (this.activeFileTasks.isEmpty()) { + List 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 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 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); } /**