diff --git a/Core/src/org/sleuthkit/autopsy/ingest/IngestManager.java b/Core/src/org/sleuthkit/autopsy/ingest/IngestManager.java index 76df43cae3..afee0367b0 100644 --- a/Core/src/org/sleuthkit/autopsy/ingest/IngestManager.java +++ b/Core/src/org/sleuthkit/autopsy/ingest/IngestManager.java @@ -21,6 +21,7 @@ package org.sleuthkit.autopsy.ingest; import java.beans.PropertyChangeListener; import java.beans.PropertyChangeSupport; import java.util.ArrayList; +import java.util.Collection; import java.util.HashMap; import java.util.List; import java.util.concurrent.CancellationException; @@ -61,8 +62,8 @@ public class IngestManager { private final ExecutorService startIngestJobsExecutor = Executors.newSingleThreadExecutor(); private final ExecutorService dataSourceIngestTasksExecutor = Executors.newSingleThreadExecutor(); private final ExecutorService fileIngestTasksExecutor = Executors.newFixedThreadPool(MAX_NUMBER_OF_FILE_INGEST_THREADS); - private final HashMap ingestJobs = new HashMap<>(); // Maps job ids to jobs - private final HashMap> ingestTasks = new HashMap<>(); // Maps task ids to task cancellation handles + private final HashMap ingestJobs = new HashMap<>(); // Maps job ids to jobs. Guarded by ingestJobs. + private final HashMap> ingestTasks = new HashMap<>(); // Maps task ids to task cancellation handles. Guarded by this. private AtomicLong ingestJobId = new AtomicLong(0L); private AtomicLong ingestTaskId = new AtomicLong(0L); private volatile IngestUI ingestMessageBox; @@ -124,12 +125,17 @@ public class IngestManager { * * @return True if any ingest jobs are in progress, false otherwise */ - public synchronized boolean isIngestRunning() { - return (ingestJobs.isEmpty() == false); + public boolean isIngestRunning() { + synchronized(ingestJobs) { + return (ingestJobs.isEmpty() == false); + } } - synchronized void addFileToIngestJob(long ingestJobId, AbstractFile file) { - IngestJob job = ingestJobs.get(ingestJobId); + void addFileToIngestJob(long ingestJobId, AbstractFile file) { + IngestJob job; + synchronized(ingestJobs) { + job = ingestJobs.get(ingestJobId); + } if (job != null) { scheduler.getFileIngestScheduler().queueFile(job, file); } @@ -301,44 +307,54 @@ public class IngestManager { } } - private synchronized void stopIngestTasks() { + private void stopIngestTasks() { // First mark all of the ingest jobs as cancelled. This way the // ingest modules will know they are being shut down due to // cancellation when the cancelled run ingest module tasks release // their pipelines. - for (IngestJob job : ingestJobs.values()) { + Collection jobList; + synchronized(ingestJobs) { + jobList = ingestJobs.values(); + } + for (IngestJob job : jobList) { job.cancel(); } // Cancel the run ingest module tasks, setting the state of the threads // running them to interrupted. - for (Future task : ingestTasks.values()) { - task.cancel(true); + synchronized(this) { + for (Future task : ingestTasks.values()) { + task.cancel(true); + } } - + // Jettision the remaining data source and file ingest tasks. scheduler.getFileIngestScheduler().emptyQueues(); scheduler.getDataSourceIngestScheduler().emptyQueues(); } - synchronized void reportStartIngestJobsTaskDone(long taskId) { + private synchronized void reportStartIngestJobsTaskDone(long taskId) { ingestTasks.remove(taskId); } - synchronized void reportRunIngestModulesTaskDone(long taskId) { - ingestTasks.remove(taskId); - - List completedJobs = new ArrayList<>(); - for (IngestJob job : ingestJobs.values()) { - job.releaseIngestPipelinesForThread(taskId); - if (job.areIngestPipelinesShutDown() == true) { - completedJobs.add(job.getId()); - } + private void reportRunIngestModulesTaskDone(long taskId) { + synchronized(this) { + ingestTasks.remove(taskId); } - for (Long jobId : completedJobs) { - IngestJob job = ingestJobs.remove(jobId); - fireIngestJobEvent(job.isCancelled() ? IngestEvent.INGEST_JOB_CANCELLED.toString() : IngestEvent.INGEST_JOB_COMPLETED.toString(), jobId); + List completedJobs = new ArrayList<>(); + synchronized(ingestJobs) { + for (IngestJob job : ingestJobs.values()) { + job.releaseIngestPipelinesForThread(taskId); + if (job.areIngestPipelinesShutDown() == true) { + completedJobs.add(job.getId()); + } + } + + for (Long jobId : completedJobs) { + IngestJob job = ingestJobs.remove(jobId); + fireIngestJobEvent(job.isCancelled() ? IngestEvent.INGEST_JOB_CANCELLED.toString() : IngestEvent.INGEST_JOB_COMPLETED.toString(), jobId); + } } } @@ -384,7 +400,7 @@ public class IngestManager { // Create an ingest job. IngestJob ingestJob = new IngestJob(IngestManager.this.ingestJobId.incrementAndGet(), dataSource, moduleTemplates, processUnallocatedSpace); - synchronized (IngestManager.this) { + synchronized (ingestJobs) { ingestJobs.put(ingestJob.getId(), ingestJob); } @@ -420,7 +436,7 @@ public class IngestManager { "IngestManager.StartIngestJobsTask.run.startupErr.dlgTitle"), JOptionPane.ERROR_MESSAGE); // Jettison the ingest job and move on to the next one. - synchronized (IngestManager.this) { + synchronized (ingestJobs) { ingestJob.cancel(); ingestJobs.remove(ingestJob.getId()); }